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 docommon.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 oapache-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.
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 deextra 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:
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:
Resultados de múltiplas instruções
O plugin enviava ao XCom o resultado da última instrução. OSQLExecuteQueryOperator 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_pullsubsequente para pegar o último elemento. - Passar as instruções como uma única string separada por
;e definirsplit_statements=True. Com o valor padrãoreturn_last=True, o operator passa a enviar apenas as linhas da última instrução, igual ao plugin.
ClickHouseOperator que inseria uma lista de linhas por meio de parameters passa a ser um
SQLInsertRowsOperator:
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:
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:
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.sqle remova o prefixoClickHousedo nome da classe. - Passe
conn_idexplicitamente. Os wrappers tratavam umconn_idausente ouNonecomoclickhouse_default; as classes decommon.sqlnão têm valor padrão e falham sem ele.default_args={"conn_id": "clickhouse_default"}cobre uma DAG inteira. database=nos operators ehook_params={"schema": ...}no sensor continuam funcionando; o hook do provider trataschemacomo um alias dedatabase.ClickHouseDbApiHookpassa a serClickHouseHook. Seu argumento de construtorschemaainda é 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 oclickhouse-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
- O Airflow está na versão 2.11 ou mais recente.
- Porta HTTP acessível a partir dos workers; certificados TLS válidos para o endpoint HTTP.
- Todas as conexões ClickHouse: tipo
clickhouse, porta8123ou8443, extras convertidos conforme o passo 2, verificados comairflow connections test. - Imports substituídos conforme o passo 3; prefixos
ClickHouseremovidos dos wrappers decommon.sql. clickhouse_conn_idrenomeado paraconn_idem operators e sensors, econn_iddefinido em toda task que dependia do padrão do plugin.settings=movido para dentro dehook_params={"session_settings": ...}.hook.execute("INSERT ... VALUES", rows)substituído porbulk_insert_rows;ClickHouseOperator(parameters=rows)substituído porSQLInsertRowsOperator.- Callables de sensor ajustados de
result[0][0]para o valor puro da cell. - Usos de
with_column_types,external_tables,columnar,query_idetypes_checkreescritos com umhandlerouget_client(). - Consumers downstream de XComs vindos de tasks com múltiplas instruções e de INSERT revisados.
- Código que captura exceptions do
clickhouse_driveratualizado. airflow-clickhouse-plugineclickhouse-driverdesinstalados.