Skip to main content
До появления официального провайдера ClickHouse большинство пользователей Airflow подключались к ClickHouse через пакет сообщества airflow-clickhouse-plugin. В этом руководстве описан перенос существующего развертывания на apache-airflow-providers-clickhousedb. О том, как работает сам провайдер, см. Подключение Apache Airflow к ClickHouse.

Зачем нужен нативный провайдер?

Провайдеры — это механизм интеграции Airflow со сторонними системами. Они выпускаются, тестируются и документируются вместе со всей остальной экосистемой Airflow, а данный провайдер поддерживается сообществом Airflow совместно с командой ClickHouse. Он построен на clickhouse-connect — клиенте Python, который разрабатывает и поддерживает сама компания ClickHouse, а не на драйвере, сопровождаемом сообществом. Благодаря этому новые возможности сервера и исправления доходят до пользователей Airflow по поддерживаемому пути. Переход на провайдер даёт вам пакет с официальным местом размещения, стандартные операторы и сенсоры common.sql, а также тип соединения, который отображается в интерфейсе Airflow так же, как и для любой другой базы данных. Эти два пакета различаются не только путями импорта. Плагин обращается к ClickHouse по собственному TCP-протоколу с помощью clickhouse-driver. Провайдер работает по HTTP(S) через clickhouse-connect и использует универсальные операторы common.sql вместо того, чтобы поставлять собственные, специфичные для ClickHouse.
Прочитайте руководство целиком, прежде чем что-либо менять. В частности, изменение соединения затрагивает сразу все DAG.

Краткий обзор

Шаг 1: Проверьте prerequisites и выполните установку

Для работы provider требуется Airflow 2.11 или новее, а также apache-airflow-providers-common-sql 1.32.0 или новее. Если у вас более старый release, сначала обновите Airflow.
Эти два пакета находятся в разных пространствах имен Python, поэтому их можно установить одновременно и переносить DAG по одному. Удалите плагин, когда его больше ничто не импортирует:

Шаг 2. Обновите соединения

Именно на этом шаге всё ломается, если его пропустить. Все существующие соединения с ClickHouse указывают на нативный порт, а провайдеру нужен HTTP-порт. Если ClickHouse находится за firewall или load balancer, перед переключением убедитесь, что HTTP-порт доступен с воркеров. Проверьте, что HTTP interface включён на сервере (http_port или https_port в конфигурации сервера). ClickHouse Cloud предоставляет HTTPS только на порту 8443. Соединения, сохранённые в виде URI (clickhouse://user:pass@host:9000/db?secure=true), уже имеют тип соединения clickhouse, поскольку Airflow определяет его по scheme URI. Для них меняются только порт и дополнительные ключи; значения из строки запроса разбираются как JSON, поэтому secure=true остаётся булевым значением.

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

Плагин передавал каждый ключ из extra напрямую в clickhouse_driver.Client, поэтому соединения могут содержать любой именованный аргумент clickhouse-driver. Провайдер читает только фиксированный набор ключей, а всё остальное передаёт через client_kwargs. Соответствие следующее: До:
После:
Проверьте каждое перенесённое соединение, прежде чем вносить изменения в DAG:

Шаг 3: Замените импорты

Обёртки common.sql с префиксом ClickHouse были нужны лишь для того, чтобы подставить hook плагина. Provider регистрирует тип соединения clickhouse, поэтому классы common.sql без префикса сами получают hook из соединения. Если вы использовали эти обёртки, миграция обычно сводится к правке строки импорта, удалению префикса ClickHouse и явной передаче conn_id (см. шаг 7).

Шаг 4: от ClickHouseOperator к SQLExecuteQueryOperator

До:
После:

Результаты нескольких операторов

Плагин отправлял в XCom результат последнего оператора. SQLExecuteQueryOperator возвращает по одному результату на каждый оператор, если sql задан списком, поэтому приведённый выше пример отправляет [[], [(12345.0,)]] вместо [(12345.0,)], которые отправлял плагин. Выберите один из вариантов:
  • Изменить нижестоящий xcom_pull так, чтобы он брал последний элемент.
  • Передать операторы одной строкой, разделив их символом ;, и задать split_statements=True. Тогда при значении по умолчанию return_last=True оператор отправит только строки последнего оператора — как и плагин.
ClickHouseOperator, который вставлял список строк через parameters, становится SQLInsertRowsOperator:
Всегда передавайте columns; без него оператор будет искать таблицу через SQLAlchemy.

Сохранение типов столбцов

with_column_types=True возвращал (rows, [(name, type), ...]). Это поведение можно воспроизвести с помощью handler; курсор clickhouse-connect возвращает имена типов ClickHouse в cursor.description:

Шаг 5: от ClickHouseHook.execute к методам DbApiHook

Hook плагина предоставлял единственный метод — execute, повторяющий clickhouse_driver.Client.execute. Hook провайдера является DbApiHook, поэтому он получает стандартные методы, доступные у любого другого SQL-провайдера: run, get_records, get_first, get_pandas_df, get_df, insert_rows и test_connection. Аргументы конструктора clickhouse_conn_id и database остались без изменений. fetch_all_handler и остальные handlers импортируются из airflow.providers.common.sql.hooks.handlers. Самый распространённый приём работы с hook в коде эпохи плагина — массовая вставка. Обычным вызовом run её не выполнить: курсор DB-API попытается подставить строки в SQL-строку. Вместо этого используйте нативную вставку: До:
После:
bulk_insert_rows требует column_names. batch_size необязателен и ограничивает потребление памяти при очень больших объёмах входных данных. Универсальный метод insert_rows(table, rows, target_fields=[...], executemany=True) также завершается нативной вставкой, но только при executemany=True; по умолчанию на каждую строку отправляется отдельный HTTP-запрос. Для всего, что не покрывается интерфейсом DB-API, get_client() возвращает необработанный клиент clickhouse-connect, настроенный на основе соединения Airflow. Он заменяет собой все специфичные для clickhouse-driver аргументы, которые предоставлял плагин:

Шаг 6: с ClickHouseSensor на SqlSensor

Это единственная замена, при которой меняются входные данные вызываемого объекта. Поскольку плагин передавал весь результирующий набор, рабочий код сенсора обращается к нему по индексу. При миграции уберите такое индексирование: До:
После:
Если вызываемому объекту нужна вся строка, передайте selector=lambda row: row. Если нужны все строки, реализуйте проверку на SQL так, чтобы запрос возвращал одно булевое значение или количество.

Шаг 7: семейство обёрток common.sql

Код, использовавший ClickHouseSQLExecuteQueryOperator, ClickHouseSqlSensor и другие обёртки с префиксом ClickHouse, потребует минимальных изменений:
  • Измените импорт на модуль common.sql и уберите префикс ClickHouse из имени класса.
  • Передавайте conn_id явно. Обёртки воспринимали отсутствующий или равный None conn_id как clickhouse_default; у классов common.sql значения по умолчанию нет, и без него они завершаются с ошибкой. Параметр default_args={"conn_id": "clickhouse_default"} закрывает весь DAG.
  • database= у операторов и hook_params={"schema": ...} у сенсора продолжают работать: hook провайдера воспринимает schema как алиас database.
  • ClickHouseDbApiHook становится ClickHouseHook. Аргумент конструктора schema по-прежнему принимается как алиас; database — нативное написание для ClickHouse, и при указании обоих приоритет остаётся за ним.
  • Соединению по-прежнему необходимы изменения порта и extras из шага 2 — обёртки тоже использовали собственный протокол.

Различия в поведении, требующие проверки

Даже после успешной компиляции кода некоторые вещи во время выполнения ведут себя иначе. Значение XCom для задач INSERT. Плагин передавал то, что возвращал clickhouse-driver, а для вставки VALUES с параметрами это было количество вставленных строк. Провайдер передаёт пустой результирующий набор для операторов, которые не возвращают строк. Нижестоящие задачи, читающие количество строк из XCom, должны получать его иначе — например, с помощью последующего SELECT count(). Настройки сеанса и операторы SET. Оба пакета выполняют список из нескольких операторов по одному соединению, а clickhouse-connect по умолчанию создаёт отдельный сеанс для каждого клиента, поэтому оператор SET в начале списка всё равно должен применяться к последующим операторам. И всё же предпочтительнее использовать session_settings: этот способ явный, поддерживает шаблонизацию и работает одинаково независимо от того, сохраняется ли сеанс на сервере или в промежуточном proxy. Если ваши DAG зависят от SET, проверьте поведение в своей среде. Исключения. Теперь ошибки представлены классами clickhouse_connect.driver.exceptions.DatabaseError, OperationalError или ProgrammingError вместо clickhouse_driver.errors.ServerException и NetworkError. Обновите блоки except, код on_failure_callback и логику повторных попыток, которая анализирует типы исключений. Сжатие. clickhouse-driver не включал сжатие, если не был задан параметр compression; clickhouse-connect по умолчанию включает сжатие HTTP-ответов и согласует алгоритм с сервером. Чтобы вернуть прежнее поведение, укажите "compress": false в дополнительных параметрах соединения. Пакет clickhouse-cityhash, который требовался плагину для сжатия по собственному протоколу, больше не нужен; lz4 остаётся установленным как зависимость clickhouse-connect. Сопоставление типов. Оба драйвера возвращают нативные типы Python, но это разные кодовые базы. Проверьте задачи, зависящие от точных типов для DateTime64 с часовыми поясами, Decimal, UUID, столбцов Nullable и вложенных значений Array или Map, особенно там, где результат передаётся в XCom и используется далее по цепочке. Идентификация запросов в system.query_log. Запросы теперь поступают через HTTP interface, поэтому отображаются со значением interface = 2 вместо 1, а столбец http_user_agent таблицы system.query_log содержит версии Airflow и провайдера, а также дополнительный параметр client_name, если он задан. Любой мониторинг, фильтровавший по собственному протоколу или по имени клиента clickhouse-driver, потребует обновления. Запрос SELECT, не возвращающий строк, создаёт вторую запись: курсор DB-API выполняет SELECT * FROM (...) LIMIT 0, чтобы получить метаданные столбцов. Тайм-ауты по HTTP. send_receive_timeout теперь соответствует тайм-ауту чтения HTTP, а любой proxy или load balancer между воркерами и ClickHouse применяет к запросу собственный idle timeout. Для операторов, которые выполнялись много минут по собственному протоколу, эти лимиты может потребоваться увеличить. Работа с соединениями. Hook создаёт клиент clickhouse-connect на каждый вызов run или get_client — так же, как плагин открывал новое native connection на каждый execute. Клиенты используют общий HTTP connection pool в рамках процесса, поэтому вызов close() для клиента, полученного через get_client(), — скорее хорошая практика, чем необходимость; пул освобождается при завершении процесса задачи. Client является менеджером контекста, поэтому with hook.get_client() as client: — самая аккуратная форма.

Контрольный список

  1. Airflow версии 2.11 или новее.
  2. HTTP-порт доступен с воркеров; TLS-сертификаты действительны для HTTP-конечной точки.
  3. Для каждого соединения с ClickHouse: тип clickhouse, порт 8123 или 8443, extras преобразованы согласно шагу 2, проверка выполнена командой airflow connections test.
  4. Импорты заменены согласно шагу 3; префиксы ClickHouse убраны из обёрток common.sql.
  5. clickhouse_conn_id переименован в conn_id у операторов и сенсоров, а conn_id задан для каждой задачи, которая полагалась на значение по умолчанию из плагина.
  6. settings= перенесены в hook_params={"session_settings": ...}.
  7. hook.execute("INSERT ... VALUES", rows) заменён на bulk_insert_rows; ClickHouseOperator(parameters=rows) заменён на SQLInsertRowsOperator.
  8. Вызываемые объекты сенсоров переведены с result[0][0] на непосредственное значение ячейки.
  9. Случаи использования with_column_types, external_tables, columnar, query_id и types_check переписаны с применением handler или get_client().
  10. Проверены нижестоящие потребители XCom из задач с несколькими операторами и задач INSERT.
  11. Код, перехватывающий исключения clickhouse_driver, обновлён.
  12. airflow-clickhouse-plugin и clickhouse-driver удалены.

Использование ИИ-ассистента для написания кода

Приведённое выше сопоставление намеренно сделано механическим, чтобы ИИ-ассистент мог применить его к репозиторию с DAG. Промпт, который хорошо себя показал:
Просмотрите diff. Изменения в соединениях, семантика сенсоров и потребители XCom — именно те места, где автоматический рефакторинг чаще всего даёт сбой.
Последнее изменение 26 сентября 2026 г.