Apache Airflow é uma plataforma de código aberto para criar, agendar e monitorar fluxos de trabalho como código. Os fluxos de trabalho são definidos como grafos acíclicos direcionados (DAGs) de tarefas escritas em Python.
O provedor apache-airflow-providers-clickhousedb conecta o Airflow ao ClickHouse, permitindo executar consultas, criar tabelas e carregar dados como parte de um DAG. Ele se conecta pela interface HTTP usando o cliente clickhouse-connect e expõe o ClickHouse por meio do framework SQL padrão do Airflow, para que o SQLExecuteQueryOperator padrão lide com DDL, DML e consultas analíticas sem exigir um operador específico do ClickHouse.
Instale o provedor
Instale o provedor no ambiente onde o scheduler e os workers do Airflow são executados:
pip install apache-airflow-providers-clickhousedbO provedor depende de apache-airflow-providers-common-sql e clickhouse-connect, que são instalados junto com ele. Para passar os resultados da consulta para DataFrames do pandas ou do polars, instale os extras opcionais:
pip install 'apache-airflow-providers-common-sql[pandas,polars]'Criar uma conexão com o ClickHouse
O provedor registra um tipo de conexão clickhouse. Crie uma conexão pela UI do Airflow em Admin > Connections ou defina uma pela CLI ou por uma variável de ambiente.
Na UI, selecione ClickHouse como tipo de conexão e preencha os campos:
| Campo | Descrição | Padrão |
|---|---|---|
| Host | Hostname do servidor ClickHouse, por exemplo abc123.clickhouse.cloud |
localhost |
| Port | Porta HTTP(S) | 8123 (sem TLS), 8443 (TLS) |
| Login | Nome de usuário do ClickHouse | default |
| Password | Senha do usuário do ClickHouse | (vazio) |
| Database | Banco de dados padrão da conexão. Na UI, esse campo aparece como Database; ao definir a conexão por URI ou JSON, ele corresponde ao campo schema. |
default |
Para o ClickHouse Cloud ou qualquer cluster self-hosted com TLS habilitado, defina secure como true no campo Extra e use a porta TLS (8443).
Opções extras de conexão
O provedor disponibiliza opções adicionais como campos específicos no formulário de conexão. Se, em vez disso, você definir a conexão por URI, JSON ou variável de ambiente, informe essas opções como chaves no objeto JSON extra. Todas são opcionais:
chave extra |
campo da UI | padrão | descrição |
|---|---|---|---|
secure |
Usar TLS (HTTPS) | false |
Habilita HTTPS/TLS. |
verify |
Verificar certificado SSL | true |
Verifica o certificado TLS do servidor quando secure é true. Defina false para certificados autoassinados. |
connect_timeout |
Tempo limite de conexão (segundos) | 10 |
Tempo limite da conexão HTTP, em segundos. |
send_receive_timeout |
Tempo limite da consulta (segundos) | 300 |
Tempo limite de leitura/gravação da consulta, em segundos. Aumente esse valor para consultas analíticas de longa duração. |
compress |
Habilitar compactação LZ4 | true |
Habilita a compactação LZ4 dos resultados. |
client_name |
Nome do cliente | (vazio) | Um rótulo acrescentado ao identificador da versão do Airflow no User-Agent do ClickHouse e à coluna client_name de system.query_log. |
session_settings |
Configurações da sessão (JSON) | (vazio) | Configurações de sessão do ClickHouse aplicadas a cada consulta nessa conexão, por exemplo {"max_execution_time": 300, "max_threads": 8}. |
client_kwargs |
kwargs do cliente (JSON) | (vazio) | Argumentos de palavra-chave adicionais encaminhados para clickhouse_connect.get_client(), por exemplo um http_proxy. |
Defina uma conexão sem a UI
Defina a conexão por meio de uma variável de ambiente. O formato de URI abrange host, credenciais e banco de dados:
export AIRFLOW_CONN_CLICKHOUSE_DEFAULT='clickhouse://default:password@localhost:8123/my_database'Todos os componentes do URI devem ser codificados para URL. Para TLS, timeouts ou configurações de sessão, use o formato JSON, que expõe os campos Extra:
export AIRFLOW_CONN_CLICKHOUSE_DEFAULT='{
"conn_type": "clickhouse",
"host": "abc123.clickhouse.cloud",
"port": 8443,
"login": "default",
"password": "secret",
"schema": "my_database",
"extra": {
"secure": true,
"session_settings": {
"max_execution_time": 300,
"max_memory_usage": 10000000000
}
}
}'Todos os hooks e operadores usam o ID de conexão clickhouse_default, a menos que você especifique outro.
Executar consultas com SQLExecuteQueryOperator
Defina o conn_id do operador para a sua conexão do ClickHouse. O DAG a seguir cria uma tabela, insere linhas, lê essas linhas novamente e exclui a tabela:
from datetime import datetime
from airflow import DAG
from airflow.providers.common.sql.hooks.sql import fetch_all_handler
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
CLICKHOUSE_CONN_ID = "clickhouse_default"
CLICKHOUSE_TABLE = "airflow_example"
with DAG(
dag_id="example_clickhouse",
start_date=datetime(2021, 1, 1),
default_args={"conn_id": CLICKHOUSE_CONN_ID},
schedule="@once",
catchup=False,
) as dag:
create_table = SQLExecuteQueryOperator(
task_id="create_table",
sql=f"""
CREATE TABLE IF NOT EXISTS {CLICKHOUSE_TABLE} (
id UInt32,
name String,
ts DateTime DEFAULT now()
) ENGINE = MergeTree()
ORDER BY id
""",
)
insert_rows = SQLExecuteQueryOperator(
task_id="insert_rows",
sql=f"""
INSERT INTO {CLICKHOUSE_TABLE} (id, name) VALUES
(1, 'Alice'),
(2, 'Bob'),
(3, 'Charlie')
""",
)
read_rows = SQLExecuteQueryOperator(
task_id="read_rows",
sql=f"SELECT id, name FROM {CLICKHOUSE_TABLE} ORDER BY id",
handler=fetch_all_handler,
)
drop_table = SQLExecuteQueryOperator(
task_id="drop_table",
sql=f"DROP TABLE IF EXISTS {CLICKHOUSE_TABLE}",
)
create_table >> insert_rows >> read_rows >> drop_tableOs resultados da consulta são recuperados com o handler padrão (fetch_all_handler). Para retornar algo diferente do conjunto completo de resultados, passe um handler diferente, como fetch_one_handler, para retornar apenas a primeira linha.
Use um banco de dados diferente por tarefa
Quando uma conexão aponta para um cluster e tarefas individuais fazem consultas em bancos de dados diferentes, sobrescreva o banco de dados por meio de hook_params em vez de criar uma conexão separada:
read_rows = SQLExecuteQueryOperator(
task_id="read_rows",
conn_id=CLICKHOUSE_CONN_ID,
sql="SELECT count() FROM events",
hook_params={"database": "analytics"},
)Use o hook diretamente
Para casos que não se encaixam em um operador SQL — inserção em massa, streaming ou chamadas específicas do cliente ClickHouse — use ClickHouseHook dentro de uma tarefa em Python.
O método bulk_insert_rows do hook usa o caminho nativo de inserção colunar em clickhouse-connect, que é muito mais rápido do que inserções linha por linha para grandes volumes de dados. Defina batch_size para limitar o pico de memória em entradas muito grandes:
from airflow.providers.clickhousedb.hooks.clickhouse import ClickHouseHook
hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default")
hook.bulk_insert_rows(
table="events",
rows=[("user1", "click"), ("user2", "view")],
column_names=["user_id", "action"],
batch_size=1000,
)Chame get_client() para acessar o client subjacente do clickhouse-connect para tudo o que o hook não expõe diretamente:
client = hook.get_client()
total = client.query("SELECT count() FROM events").result_rows[0][0]Aplicar configurações de sessão
Passe configurações de sessão ao construir o hook, seja diretamente ou por meio do hook_params de um operador. As configurações passadas ao construtor são mescladas com quaisquer session_settings definidas no campo Extra da conexão, e os valores do construtor prevalecem em caso de conflito entre chaves:
hook = ClickHouseHook(
clickhouse_conn_id="clickhouse_default",
session_settings={"max_execution_time": 60, "max_threads": 4},
)