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