Skip to main content
Antes de o provider oficial do ClickHouse existir, a maioria dos usuários do Airflow se conectava ao ClickHouse por meio do pacote da comunidade airflow-clickhouse-plugin. Este guia mostra como migrar uma implantação existente para o apache-airflow-providers-clickhousedb. Para saber como o próprio provider funciona, consulte Conectar o Apache Airflow ao ClickHouse.

Por que um provider nativo?

Os providers são a forma pela qual o Airflow se integra a sistemas de terceiros. Eles são lançados, testados e documentados junto com o restante do ecossistema Airflow, e este em particular é mantido pela comunidade Airflow em conjunto com a equipe da ClickHouse. Ele é construído sobre o clickhouse-connect, o cliente Python que a própria ClickHouse desenvolve e mantém, e não sobre um driver mantido pela comunidade. Dessa forma, novos recursos e correções do servidor chegam aos usuários do Airflow por um caminho com suporte oficial. Migrar para o provider oferece um package com uma casa oficial, os operators e sensors padrão do common.sql e um tipo de connection que aparece na UI do Airflow como o de qualquer outro banco de dados. A diferença entre os dois packages vai além dos caminhos de import. O plugin se comunica com o ClickHouse pelo protocolo TCP nativo, usando o clickhouse-driver. Já o provider se comunica por HTTP(S), usando o clickhouse-connect, e se integra aos operators genéricos do common.sql em vez de fornecer operators específicos para o ClickHouse.
Leia o guia inteiro antes de alterar qualquer coisa. A mudança na connection, em especial, afeta todas as DAGs ao mesmo tempo.

Visão geral

Passo 1: Verificar os prerequisites e instalar

O provider exige o Airflow 2.11 ou mais recente e o apache-airflow-providers-common-sql 1.32.0 ou mais recente. Faça primeiro o upgrade do Airflow caso você esteja em um lançamento mais antigo.
Os dois pacotes ficam em espaços de nomes distintos do Python, portanto podem ser instalados lado a lado enquanto você migra DAG por DAG. Remova o plugin quando nada mais o importar:

Passo 2: Atualizar as conexões

Este é o passo que quebra tudo se você pulá-lo. Toda conexão existente do ClickHouse aponta para a porta native, e o provider precisa da porta HTTP. Se o ClickHouse estiver atrás de um firewall ou de um balanceador de carga, verifique se a porta HTTP está acessível a partir dos workers antes de fazer a troca. Confirme se a interface HTTP está habilitada no servidor (http_port ou https_port na configuração do servidor). O ClickHouse Cloud expõe HTTPS apenas na porta 8443. Conexões armazenadas como URIs (clickhouse://user:pass@host:9000/db?secure=true) já têm o tipo de connection clickhouse, pois o Airflow o deriva do scheme da URI. Nesses casos, mudam apenas a porta e as chaves extra; os valores da query string são interpretados como JSON, portanto secure=true permanece um Boolean.

Extras de conexão

O plugin repassava todas as chaves de extra diretamente para clickhouse_driver.Client, portanto as conexões podem carregar qualquer keyword argument do clickhouse-driver. O provider lê apenas um conjunto fixo de chaves e encaminha as demais por meio de client_kwargs. Faça a conversão da seguinte forma: Antes:
Depois:
Verifique cada connection migrada antes de alterar as DAGs:

Passo 3: substituir os imports

Os wrappers de common.sql com prefixo ClickHouse existiam apenas para injetar o hook do plugin. O provider registra o tipo de connection clickhouse, de modo que as classes common.sql sem prefixo resolvem o hook a partir da connection por conta própria. Se você usava esses wrappers, a migration geralmente se resume a ajustar a linha de import, remover o prefixo ClickHouse e passar conn_id explicitamente (consulte o passo 7).

Passo 4: ClickHouseOperator para SQLExecuteQueryOperator

Antes:
Depois:

Resultados de múltiplas instruções

O plugin enviava ao XCom o resultado da última instrução. O SQLExecuteQueryOperator retorna um resultado por instrução quando sql é uma lista, portanto o exemplo acima envia [[], [(12345.0,)]] onde o plugin enviava [(12345.0,)]. Escolha uma destas opções:
  • Alterar o xcom_pull subsequente para pegar o último elemento.
  • Passar as instruções como uma única string separada por ; e definir split_statements=True. Com o valor padrão return_last=True, o operator passa a enviar apenas as linhas da última instrução, igual ao plugin.
Um ClickHouseOperator que inseria uma lista de linhas por meio de parameters passa a ser um SQLInsertRowsOperator:
Sempre passe columns; sem isso, o operator busca a tabela via SQLAlchemy.

Preservando os tipos de coluna

with_column_types=True retornava (rows, [(name, type), ...]). Reproduza esse comportamento com um handler; o cursor do clickhouse-connect informa os nomes dos tipos do ClickHouse em cursor.description:

Passo 5: de ClickHouseHook.execute para os métodos do DbApiHook

O hook do plugin expunha um único método, execute, espelhando clickhouse_driver.Client.execute. O hook do provider é um DbApiHook, portanto ganha os métodos padrão que todos os outros providers SQL têm: run, get_records, get_first, get_pandas_df, get_df, insert_rows e test_connection. Os argumentos de construtor clickhouse_conn_id e database permanecem inalterados. O fetch_all_handler e os demais handlers são importados de airflow.providers.common.sql.hooks.handlers. O uso mais comum do hook em código da era dos plugins é o bulk insert. Ele não pode ser uma simples chamada a run, porque o cursor DB-API tentaria formatar as linhas dentro da string SQL. Em vez disso, use o native insert: Antes:
Depois:
bulk_insert_rows exige column_names. batch_size é opcional e limita o uso de memória em entradas muito grandes. O insert_rows(table, rows, target_fields=[...], executemany=True) genérico também resulta em um insert nativo, mas apenas com executemany=True; por padrão, é enviada uma requisição HTTP por linha. Para tudo o que a superfície DB-API não cobre, get_client() retorna o client clickhouse-connect bruto, configurado a partir da connection do Airflow. Ele substitui cada argument específico do clickhouse-driver que o plugin expunha:

Etapa 6: ClickHouseSensor para SqlSensor

Esta é a única substituição em que a entrada do callable muda. Como o plugin entregava o conjunto de resultados completo, o código de sensor existente faz indexação sobre ele. Remova essa indexação ao migrar: Antes:
Depois:
Se o seu callable precisar da linha inteira, passe selector=lambda row: row. Se ele precisar de todas as linhas, escreva a verificação em SQL para que a consulta retorne um único booleano ou uma contagem.

Passo 7: A família de wrappers common.sql

O código que usava ClickHouseSQLExecuteQueryOperator, ClickHouseSqlSensor e os demais wrappers com prefixo ClickHouse é o que exige menos trabalho:
  • Altere o import para o módulo common.sql e remova o prefixo ClickHouse do nome da classe.
  • Passe conn_id explicitamente. Os wrappers tratavam um conn_id ausente ou None como clickhouse_default; as classes de common.sql não têm valor padrão e falham sem ele. default_args={"conn_id": "clickhouse_default"} cobre uma DAG inteira.
  • database= nos operators e hook_params={"schema": ...} no sensor continuam funcionando; o hook do provider trata schema como um alias de database.
  • ClickHouseDbApiHook passa a ser ClickHouseHook. Seu argumento de construtor schema ainda é aceito como alias; database é a grafia nativa do ClickHouse e tem precedência quando ambos são informados.
  • A connection ainda precisa das alterações de porta e extras do passo 2. Os wrappers também usavam o protocolo nativo.

Diferenças de comportamento a revisar

Mesmo depois de o código compilar, algumas coisas se comportam de forma diferente em tempo de execução. Valor de XCom das tasks de INSERT. O plugin enviava o que o clickhouse-driver retornava, o que, para um insert VALUES com parâmetros, era a contagem de linhas inseridas. O provider envia um conjunto de resultados vazio para instruções que não retornam linhas. Tasks subsequentes que leem a contagem de linhas do XCom precisam obtê-la de outra forma, por exemplo com um SELECT count() posterior. Configurações de sessão versus instruções SET. Ambos os pacotes executam uma lista de múltiplas instruções em uma única connection, e o clickhouse-connect cria uma session por client por padrão, de modo que uma instrução SET no início da lista ainda deve valer para as instruções posteriores. Ainda assim, prefira session_settings: é explícito e templatizado, e funciona da mesma maneira, mantendo o servidor ou um proxy intermediário a session ou não. Confirme o comportamento no seu ambiente se suas DAGs dependem de SET. Exceções. Os erros agora são clickhouse_connect.driver.exceptions.DatabaseError, OperationalError ou ProgrammingError em vez de clickhouse_driver.errors.ServerException e NetworkError. Atualize as cláusulas except, o código de on_failure_callback e a lógica de retentativa que inspeciona tipos de exceção. Compressão. O clickhouse-driver deixava a compressão desativada a menos que compression fosse definido; o clickhouse-connect habilita a compressão da resposta HTTP por padrão e negocia o algoritmo com o servidor. Defina "compress": false no extra da connection para restaurar o comportamento antigo. O pacote clickhouse-cityhash, de que o plugin precisava para compressão nativa, não é mais necessário; o lz4 continua instalado como dependência do clickhouse-connect. Mapeamento de tipos. Ambos os drivers retornam tipos nativos do Python, mas são bases de código diferentes. Revise as tasks que dependem de tipos exatos para DateTime64 com fusos horários, Decimal, UUID, colunas Nullable e valores Array ou Map aninhados, especialmente quando o resultado é enviado ao XCom e consumido por tasks subsequentes. Identificação de consultas em system.query_log. As consultas agora chegam pela interface HTTP, então aparecem com interface = 2 em vez de 1, e a coluna http_user_agent de system.query_log carrega as versões do Airflow e do provider, além do extra client_name, se definido. Qualquer monitoramento que filtrava pelo protocolo nativo ou pelo nome de client do clickhouse-driver precisa ser atualizado. Um SELECT que não retorna linhas gera uma segunda entrada: o cursor DB-API executa SELECT * FROM (...) LIMIT 0 para recuperar os metadados das colunas. Tempos limite sobre HTTP. O send_receive_timeout agora é o tempo limite de leitura HTTP, e qualquer proxy ou balanceador de carga entre os workers e o ClickHouse aplica seu próprio idle timeout à requisição. Instruções que levavam muitos minutos pelo protocolo nativo podem exigir que esses limites sejam aumentados. Tratamento de connections. O hook cria um client clickhouse-connect a cada chamada de run ou get_client, espelhando a forma como o plugin abria uma nova connection nativa a cada execute. Os clients compartilham um HTTP connection pool no escopo do processo, então chamar close() em um client obtido de get_client() é boa prática, não uma obrigação; o pool é liberado quando o processo da task termina. O client é um context manager, portanto with hook.get_client() as client: é a forma mais elegante.

Checklist

  1. O Airflow está na versão 2.11 ou mais recente.
  2. Porta HTTP acessível a partir dos workers; certificados TLS válidos para o endpoint HTTP.
  3. Todas as conexões ClickHouse: tipo clickhouse, porta 8123 ou 8443, extras convertidos conforme o passo 2, verificados com airflow connections test.
  4. Imports substituídos conforme o passo 3; prefixos ClickHouse removidos dos wrappers de common.sql.
  5. clickhouse_conn_id renomeado para conn_id em operators e sensors, e conn_id definido em toda task que dependia do padrão do plugin.
  6. settings= movido para dentro de hook_params={"session_settings": ...}.
  7. hook.execute("INSERT ... VALUES", rows) substituído por bulk_insert_rows; ClickHouseOperator(parameters=rows) substituído por SQLInsertRowsOperator.
  8. Callables de sensor ajustados de result[0][0] para o valor puro da cell.
  9. Usos de with_column_types, external_tables, columnar, query_id e types_check reescritos com um handler ou get_client().
  10. Consumers downstream de XComs vindos de tasks com múltiplas instruções e de INSERT revisados.
  11. Código que captura exceptions do clickhouse_driver atualizado.
  12. airflow-clickhouse-plugin e clickhouse-driver desinstalados.

Usando um assistente de programação com IA

O mapeamento acima é deliberadamente mecânico, de modo que um assistente de programação consiga aplicá-lo a um repositório de DAGs. Um prompt que funcionou bem:
Revise o diff. As mudanças de conexão, a semântica dos sensores e os consumers do XCom são os pontos onde as reescritas automatizadas dão errado.
Última modificação em 26 de setembro de 2026