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