Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Conecte o Apache Airflow ao ClickHouse

Suportado pelo ClickHouse

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-clickhousedb

O 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_table

Os 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},
)
Navigation