Плагин Airflow для ClickHouse
Самый популярный плагин Apache Airflow для ClickHouse, один из самых скачиваемых плагинов на PyPI, построенный на базе супер драйвера mymarilyn/clickhouse-driver.
Этот плагин предоставляет два семейства операторов: расширенный clickhouse_driver.Client.execute-based и стандартный, совместимый с Python DB API 2.0.
Оба семейства операторов снабжены тестами для различных версий Airflow и Python.
ClickHouseOperator
-
ClickHouseHook -
ClickHouseSensor
Эти операторы основаны на методе и аргументах Client.execute драйвера mymarilyn/clickhouse-driver. Они предлагают полный функционал драйвера clickhouse. Если Вы только начинаете работать с ClickHouse в Airflow, обратите свое внимание именно на них.
Характеристики:
- Шаблонизация SQL: SQL-запросы и другие параметры могут быть выполнены по шаблону.
- Множественные SQL-запросы: выполнение нескольких SQL-запросов в рамках одного ClickHouseOperator. Результат последнего запроса передается в XCom (настраивается командой do_xcom_push).
- Ведение журнала: Выполненные запросы записываются в журнал в удобном для восприятия формате, что облегчает их отслеживание и отладку.
- Эффективный родной протокол ClickHouse: Использует эффективный родной TCP-протокол ClickHouse благодаря драйверу clickhouse-driver. Не поддерживает протокол HTTP.
- Пользовательские параметры соединения: Поддерживает дополнительные параметры соединения ClickHouse, такие как различные таймауты, сжатие, безопасность, через свойство Connection.extra.
Примеры представлены по ссылке.
Установка и зависимости:
pip install -U airflow-clickhouse-plugin
Зависимости: только apache-airflow и clickhouse-driver.
Семейство Python DB API 2.0
-
Операторы:
-
ClickHouseSQLExecuteQueryOperator -
ClickHouseSQLColumnCheckOperator -
ClickHouseSQLTableCheckOperator -
ClickHouseSQLCheckOperator -
ClickHouseSQLValueCheckOperator -
ClickHouseSQLIntervalCheckOperator -
ClickHouseSQLThresholdCheckOperator -
ClickHouseBranchSQLOperator
-
-
ClickHouseDbApiHook -
ClickHouseSqlSensor
Эти операторы объединяют clickhouse_driver.dbapi и apache-airflow-providers-common-sql. Хотя по сравнению с Client.execute они имеют ограниченную функциональность (поддерживаются не все аргументы), все же они предоставляют стандартизированный интерфейс, что особенно полезно при переносе конвейеров Airflow на ClickHouse с другого SQL-провайдера, поддерживающего пакет common.sql Airflow, например MySQL, Postgres, BigQuery или других.
Набор функций этой версии полностью основан на провайдере common.sql Airflow.
Пример доступен по ссылке.
Установка и зависимости
Добавьте common.sql при установке плагина: pip install -U airflow-clickhouse-plugin[common.sql] - для включения операторов DB API 2.0.
Зависимости: apache-airflow-providers-common-sql (обычно поставляется в комплекте с Airflow) в дополнение к apache-airflow и clickhouse-driver.
Поддержка версий Python и Airflow
Разные версии плагина поддерживают различные комбинации версий Python и Airflow. В основном мы поддерживаем Airflow 2.0+ и Python 3.8+. Если Вам нужно использовать плагин с более старыми комбинациями Python-Airflow, выберите подходящую версию плагина:
|
Версия плагина airflow-clickhouse-plugin |
Версия Airflow |
Версия Python |
|---|---|---|
|
1.4.0 |
>=2.0.0,<2.11.0 |
~=3.8 |
|
1.3.0 |
>=2.0.0,<2.10.0 |
~=3.8 |
|
1.2.0 |
>=2.0.0,<2.9.0 |
~=3.8 |
|
1.1.0 |
>=2.0.0,<2.8.0 |
~=3.8 |
|
1.0.0 |
>=2.0.0,<2.7.0 |
~=3.8 |
|
0.11.0 |
~=2.0.0,>=2.2.0,<2.7.0 |
~=3.7 |
|
0.10.0,0.10.1 |
~=2.0.0,>=2.2.0,<2.6.0 |
~=3.7 |
|
0.9.0,0.9.1 |
~=2.0.0,>=2.2.0,<2.5.0 |
~=3.7 |
|
0.8.2 |
>=2.0.0,<2.4.0 |
~=3.7 |
|
0.8.0,0.8.1 |
>=2.0.0,<2.3.0 |
~=3.6 |
|
0.7.0 |
>=2.0.0,<2.2.0 |
~=3.6 |
|
0.6.0 |
~=2.0.1 |
~=3.6 |
|
>=0.5.4,<0.6.0 |
~=1.10.6 |
>=2.7 or >=3.5.* |
|
>=0.5.0,<0.5.4 |
==1.10.6 |
>=2.7 or >=3.5.* |
~= означает совместимый выпуск, более подробное объяснение см. в PEP 440.
Функционал B API 2.0 подразумевает использование apache-airflow>2.9.3 (строго новее, так как версии до 2.9.3 имели ошибку, связанную с MRO, см. #87) и apache-airflow-providers-common-sql>=1.3: более ранние версии не поддерживаются.
Для предыдущих версий плагина может потребоваться дополнительный pandas: pip install airflow-clickhouse-plugin[pandas]==0.11.0. Для получения подробностей ознакомьтесь с более ранними версиями README.md.
Использование
Для запуска примеров, указанных ниже, создайте соединение Airflow и ClickHouse.
ClickHouseOperator
Для импорта ClickHouseOperator используйте команду from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator.
Аргументы:
- sql (шаблонный, обязательный): запрос (если аргумент - одна строка) или несколько запросов (итерабельность строк). Поддерживает файлы с расширением .sql.
- clickhouse_conn_id: id соединения с Airflow. Схема соединения описана ниже. id соединения, используемое по умолачнию - clickhouse_default.
- Аргументы метода clickhouse_driver.Client.execute :
- parameters (шаблонный): params метода executе. (Переименовано для того, чтобы избежать конфликта названий с аргументом params задач Airflow.)
- dict для запросов SELECT .
- list/tuple/generator для запросов INSERT .
- Если через sql передается несколько запросов, то параметры передаются всем им..
- with_column_types (не шаблонный).
- external_tables (шаблонный).
- query_id (шаблонный).
- settings (шаблонный).
- types_check (не шаблонный).
- columnar (не щаблонный).
- Документация по пргументам доступна по ссылке: clickhouse_driver.Client.execute API .
- database (шаблонный): если присутствует, переопределяет схему соединения Airflow.
- Другие аргументы (включая обязательный идентификатор task_id) наследуются от Airflow BaseOperator.
Результат последнего запроса передается в XCom (отключается с помощью аргумента do_xcom_push=False).
Другими словами, оператор является оберткой ClickHouseHook.execute method.
ClickHouseHook
Для импорта ClickHouseHook используйте коменду from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook.
kwargs конструктора (метод__init__ ):
- clickhouse_conn_id: id соединения сAirflow. Схема соединения описана ниже. id соединения, используемое по умолчанию, - clickhouse_default.
- database: если есть, то переопределяет схему соединения Airflow.
Определяет метод ClickHouseHook.execute, который просто оборачивает clickhouse_driver.Client.execute. Он имеет все те же аргументы, за исключением:
sql (вместо запроса execute): запрос (если аргумент - одна строка) или несколько запросов (итерабельность строк).
ClickHouseHook.execute возвращает результат самого последнего запроса.
Кроме того, хук определяет метод get_conn(), который возвращает базовый экземпляр clickhouse_driver.Client.
ClickHouseSensor
Для импорта ClickHouseSensor используйте команду from airflow_clickhouse_plugin.sensors.clickhouse import ClickHouseSensor.
Этот класс является оберткой метода ClickHouseHook.execute в Airflow sensor. Поддерживает все аргументы ClickHouseOperator , а также:
- is_success: вызываемый элемент, принимающий единственный аргумент - возвращаемое значение ClickHouseHook.execute. Если возвращаемое значение is_success истинно, то датчик работает успешно. По умолчанию вызываемая переменная имеет значение bool: т. е. если возвращаемое значение ClickHouseHook.execute истинно, то датчик работает успешно. Обычно execute - это список записей, возвращаемых запросом: таким образом, по умолчанию это falsy, если не возвращается ни одной записи.
- is_failure: вызываемый элемент, принимающий единственный аргумент - возвращаемое значение ClickHouseHook.execute. Если возвращаемое значение is_failure является истинным, датчик поднимает AirflowException. По умолчанию is_failure равно None, и проверка отказа не выполняется.
Как подключить Airflow к ClickHouse
В качестве типа нового соединения выберите SQLite или любую другую базу данных SQL. Специального типа соединения ClickHouse пока не существует, поэтому мы используем любой SQL как наиболее близкий.
Все атрибуты соединения необязательны: хост по умолчанию - localhost, а остальные параметры имеют значения по умолчанию, заданные clickhouse-driver. Если вы используете значения не по умолчанию, задайте их в соответствии со схемой соединения.
Если вы используете безопасное соединение с ClickHouse (это требует дополнительных настроек на стороне ClickHouse), установите extra в {«secure»:true}. Все дополнительные параметры соединения передаются в clickhouse_driver.Client как есть.
Схема соединения ClickHouse
clickhouse_driver.Client инициализируется атрибутами, хранящимися в атрибутах соединения:
|
Атрибут соединения Airflow |
Аргумент Client.__init__ |
|---|---|
|
host |
host |
|
port (int) |
port |
|
schema |
database |
|
login |
user |
|
password |
password |
|
extra |
**kwargs |
аргумент базы данных ClickHouseOperator, ClickHouseHook, ClickHouseSensor и других переопределяет атрибут схемы соединения Airflow.
Дополнительные аргументы
Вы можете задать нестандартные аргументы clickhouse_driver.Client такие как таймауты, сжатие, безопасность и т.д., используя атрибут Connection.extra . Атрибут должен содержать JSON-объект, который будет десереализован, а все его свойства будут переданы Клиенту как есть.
Например, если соединение Airflow содержит extra='{«secure»: true}', то Client.__init__ получит ключевой аргумент secure=True в дополнение к другим атрибутам соединения.
Cжатие
Для поддержки сжатия необходимо установить специальные пакеты. Например, для lz4:
pip3 install clickhouse-cityhash lz4
Тогда в uri соединения airflow необходимо включить параметр сжатия: extra='{«сжатие»: «lz4»}'. Дополнительную информацию о дополнительных параметрах вы можете получить из официальной документации драйвера clickhouse-driver.
URI соединения с компрессией будет выглядеть как clickhouse://login:password@host:port/?compression=lz4.
Чтобы узнать больше об управлении соединениями в Airflow, обратитесь к официальной документации.
Значения по умолчанию
Если какой-то атрибут соединения Airflow не установлен, он не передается в clickhouse_driver.Client. В таких случаях плагин использует значение по умолчанию из соответствующего аргумента clickhouse_driver.Connection. Например, значение user по умолчанию равно 'default'.
Это означает, что сам плагин не определяет никаких значений по умолчанию для соединения ClickHouse. Вы можете полностью полагаться на значения по умолчанию используемой Вами версии clickhouse-драйвера.
Единственным исключением является host: если атрибут соединения Airflow не задан, то используется 'localhost'.
Соединение по умолчанию
По умолчанию плагин использует соединение Airflow с идентификатором 'clickhouse_default'.
Примеры
Пример ClickHouseOperator
from airflow import DAG
from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator
from airflow.operators.python import PythonOperator
from airflow.utils.dates import days_ago
with DAG(
dag_id='update_income_aggregate',
start_date=days_ago(2),
) as dag:
ClickHouseOperator(
task_id='update_income_aggregate',
database='default',
sql=(
'''
INSERT INTO aggregate
SELECT eventDt, sum(price * qty) AS income FROM sales
WHERE eventDt = '{{ ds }}' GROUP BY eventDt
''', '''
OPTIMIZE TABLE aggregate ON CLUSTER {{ var.value.cluster_name }}
PARTITION toDate('{{ execution_date.format('%Y-%m-01') }}')
''', '''
SELECT sum(income) FROM aggregate
WHERE eventDt BETWEEN
'{{ execution_date.start_of('month').to_date_string() }}'
AND '{{ execution_date.end_of('month').to_date_string() }}'
''',
# result of the last query is pushed to XCom
),
# query_id is templated and allows to quickly identify query in ClickHouse logs
query_id='{{ ti.dag_id }}-{{ ti.task_id }}-{{ ti.run_id }}-{{ ti.try_number }}',
clickhouse_conn_id='clickhouse_test',
) >> PythonOperator(
task_id='print_month_income',
python_callable=lambda task_instance:
# pulling XCom value and printing it
print(task_instance.xcom_pull(task_ids='update_income_aggregate')),
)
Пример ClickHouseHook
from airflow import DAG
from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook
from airflow.providers.sqlite.hooks.sqlite import SqliteHook
from airflow.operators.python import PythonOperator
from airflow.utils.dates import days_ago
def sqlite_to_clickhouse():
sqlite_hook = SqliteHook()
ch_hook = ClickHouseHook()
records = sqlite_hook.get_records('SELECT * FROM some_sqlite_table')
ch_hook.execute('INSERT INTO some_ch_table VALUES', records)
with DAG(
dag_id='sqlite_to_clickhouse',
start_date=days_ago(2),
) as dag:
dag >> PythonOperator(
task_id='sqlite_to_clickhouse',
python_callable=sqlite_to_clickhouse,
)
Важное замечание: не пытайтесь вставить значения с помощью формы ch_hook.execute('INSERT INTO some_ch_table VALUES (1)'). clickhouse-драйвер требует, чтобы значения для INSERT-запроса предоставлялись через параметры (в силу специфики родного протокола ClickHouse).
Пример ClickHouseSensor
from airflow import DAG
from airflow_clickhouse_plugin.sensors.clickhouse import ClickHouseSensor
from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator
from airflow.utils.dates import days_ago
with DAG(
dag_id='listen_warnings',
start_date=days_ago(2),
) as dag:
dag >> ClickHouseSensor(
task_id='poke_events_count',
database='monitor',
sql="SELECT count() FROM warnings WHERE eventDate = '{{ ds }}'",
is_success=lambda cnt: cnt > 10000,
) >> ClickHouseOperator(
task_id='create_alert',
database='alerts',
sql='''
INSERT INTO events SELECT eventDate, count()
FROM monitor.warnings WHERE eventDate = '{{ ds }}'
''',
)
DB API 2.0: пример ClickHouseSqlSensor и ClickHouseSQLExecuteQueryOperator
from airflow import DAG
from airflow_clickhouse_plugin.sensors.clickhouse_dbapi import ClickHouseSqlSensor
from airflow_clickhouse_plugin.operators.clickhouse_dbapi import ClickHouseSQLExecuteQueryOperator
from airflow.utils.dates import days_ago
with DAG(
dag_id='listen_warnings',
start_date=days_ago(2),
) as dag:
dag >> ClickHouseSqlSensor(
task_id='poke_events_count',
hook_params=dict(schema='monitor'),
sql="SELECT count() FROM warnings WHERE eventDate = '{{ ds }}'",
success=lambda cnt: cnt > 10000,
conn_id=None, # required by common.sql SqlSensor; use None for default
) >> ClickHouseSQLExecuteQueryOperator(
task_id='create_alert',
database='alerts',
sql='''
INSERT INTO events SELECT eventDate, count()
FROM monitor.warnings WHERE eventDate = '{{ ds }}'
''',
)
Как проводить тестирование
Модульные тесты: python3 -m unittest discover -t tests -s unit
Интеграционные тесты требуют доступа к серверу ClickHouse. Вот как создать локальную тестовую среду с помощью Docker:
- Запустите сервер ClickHouse в локальном контейнере Docker: docker run -p 9000:9000 --ulimit nofile=262144:262144 -it clickhouse/clickhouse-server
- Запустите тесты с данными о подключении Airflow, заданными через переменную окружения: PYTHONPATH=src AIRFLOW_CONN_CLICKHOUSE_DEFAULT=clickhouse://localhost python3 -m unittest discover -t tests -s integration
- Остановите контейнер после выполнения тестов.
Запустите все (модульные и интеграционные) тесты с определенным соединением ClickHouse: PYTHONPATH=src AIRFLOW_CONN_CLICKHOUSE_DEFAULT=clickhouse://localhost python3 -m unittest discover -s tests
GitHub Actions
GitHub Action настроен в соответствии с данным проектом.
Тестирование внутри Docker
Запустите сервер ClickHouse внутри Docker: docker exec -it $(docker run --rm -d clickhouse/clickhouse-server) bash
Приведенная выше команда откроет bash внутри контейнера.
Установите зависимости в контейнер и запустите тесты (выполнить внутри контейнера):
apt-get update apt-get install -y python3 python3-pip git make git clone https://github.com/whisklabs/airflow-clickhouse-plugin.git cd airflow-clickhouse-plugin python3 -m pip install -r requirements.txt PYTHONPATH=src AIRFLOW_CONN_CLICKHOUSE_DEFAULT=clickhouse://localhost python3 -m unittest discover -s tests
Остановите работу контейнера.




