Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Подключите Apache Airflow к ClickHouse

Поддерживается в ClickHouse

Apache Airflow — это платформа с открытым исходным кодом для создания, планирования и мониторинга рабочих процессов как кода. Рабочие процессы определяются как ориентированные ациклические графы (DAG) задач, написанных на Python.

Провайдер apache-airflow-providers-clickhousedb подключает Airflow к ClickHouse, позволяя выполнять запросы, создавать таблицы и загружать данные в рамках DAG. Он подключается через HTTP-интерфейс с помощью клиента clickhouse-connect и предоставляет доступ к ClickHouse через общий SQL-фреймворк Airflow, поэтому стандартный SQLExecuteQueryOperator обрабатывает DDL, DML и аналитические запросы без необходимости в специальном операторе для ClickHouse.

Установите провайдер

Установите провайдер в окружение, где работают планировщик Airflow и воркеры:

pip install apache-airflow-providers-clickhousedb

Провайдер зависит от apache-airflow-providers-common-sql и clickhouse-connect, которые устанавливаются вместе с ним. Чтобы передавать результаты запросов в DataFrame из pandas или polars, установите дополнительные опциональные компоненты:

pip install 'apache-airflow-providers-common-sql[pandas,polars]'

Создайте подключение ClickHouse

Провайдер регистрирует тип подключения clickhouse. Создайте подключение в интерфейсе Airflow в разделе Admin > Connections или задайте его через CLI либо переменную окружения.

В интерфейсе выберите ClickHouse в качестве типа подключения и заполните следующие поля:

Поле Описание По умолчанию
Host Имя хоста сервера ClickHouse, например abc123.clickhouse.cloud localhost
Port Порт HTTP(S) 8123 (без шифрования), 8443 (TLS)
Login Имя пользователя ClickHouse default
Password Пароль пользователя ClickHouse (пусто)
Database База данных по умолчанию для этого подключения. В интерфейсе это поле называется Database; при задании подключения через URI или JSON это поле schema. default

Для ClickHouse Cloud или любого самоуправляемого кластера с включенным TLS установите для secure значение true в поле Extra и используйте TLS-порт (8443).

Дополнительные параметры соединения

Провайдер предоставляет дополнительные параметры в виде отдельных полей в форме соединения. Если вместо этого вы задаёте соединение через URI, JSON или переменную окружения, укажите их как ключи в объекте JSON extra. Все они необязательны:

ключ extra Поле в интерфейсе По умолчанию Описание
secure Использовать TLS (HTTPS) false Включает HTTPS/TLS.
verify Проверять SSL-сертификат true Проверяет TLS-сертификат сервера, когда secure имеет значение true. Для самоподписанных сертификатов установите false.
connect_timeout Тайм-аут соединения (секунды) 10 Тайм-аут HTTP-соединения в секундах.
send_receive_timeout Тайм-аут запроса (секунды) 300 Тайм-аут чтения/записи для запроса в секундах. Увеличьте его для длительных аналитических запросов.
compress Включить сжатие LZ4 true Включает сжатие результатов с помощью LZ4.
client_name Имя клиента (пусто) Метка, добавляемая к идентификатору версии Airflow в User-Agent ClickHouse и в столбец client_name таблицы system.query_log.
session_settings Настройки сеанса (JSON) (пусто) Настройки сеанса ClickHouse, применяемые ко всем запросам через это соединение, например {"max_execution_time": 300, "max_threads": 8}.
client_kwargs Параметры клиента (JSON) (пусто) Дополнительные именованные аргументы, передаваемые в clickhouse_connect.get_client(), например http_proxy.

Настройте подключение без интерфейса

Задайте подключение через переменную окружения. Формат URI включает хост, учетные данные и базу данных:

export AIRFLOW_CONN_CLICKHOUSE_DEFAULT='clickhouse://default:password@localhost:8123/my_database'

Все компоненты URI должны быть URL-кодированы. Для TLS, тайм-аутов и настроек сеанса используйте форму JSON, в которой доступны поля 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
        }
    }
}'

Все хуки и операторы используют идентификатор подключения clickhouse_default, если не указан другой.

Выполнение запросов с SQLExecuteQueryOperator

Укажите в conn_id оператора ваше подключение к ClickHouse. Следующий DAG создаёт таблицу, вставляет строки, считывает их обратно и удаляет таблицу:

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

Результаты запроса извлекаются с помощью обработчика handler, используемого по умолчанию (fetch_all_handler). Чтобы вернуть не весь результирующий набор, передайте другой обработчик, например fetch_one_handler, чтобы получить только первую строку.

Укажите разные базы данных для каждой задачи

Если одно соединение указывает на кластер, а отдельные задачи выполняют запросы к разным базам данных, переопределите базу данных через hook_params вместо создания отдельного соединения:

read_rows = SQLExecuteQueryOperator(
    task_id="read_rows",
    conn_id=CLICKHOUSE_CONN_ID,
    sql="SELECT count() FROM events",
    hook_params={"database": "analytics"},
)

Используйте хук напрямую

Для задач, которые не укладываются в возможности SQL-оператора, — например, для массовой вставки, стриминга или вызовов клиента, специфичных для ClickHouse, — используйте ClickHouseHook в Python-задаче.

Метод bulk_insert_rows этого хука использует нативный столбцовый путь вставки в clickhouse-connect, который на больших датасетах значительно быстрее, чем построчная вставка. Установите batch_size, чтобы ограничить пиковое потребление памяти при очень больших объёмах входных данных:

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,
)

Вызовите get_client(), чтобы получить доступ к клиенту clickhouse-connect на низком уровне для всего, что хук не предоставляет напрямую:

client = hook.get_client()
total = client.query("SELECT count() FROM events").result_rows[0][0]

Применить настройки сеанса

Передавайте настройки сеанса при создании хука — напрямую или через hook_params оператора. Настройки, переданные в конструктор, накладываются поверх любых session_settings, заданных в поле Extra подключения; при конфликте ключей приоритет имеют значения конструктора:

hook = ClickHouseHook(
    clickhouse_conn_id="clickhouse_default",
    session_settings={"max_execution_time": 60, "max_threads": 4},
)
Navigation