Провайдер DuckDB
Обзор
Провайдер DuckDB позволяет Airflow выполнять SQL-запросы к базам данных DuckDB непосредственно из DAG. Данный провайдер не включает среду выполнения DuckDB — вместо этого предполагается, что на каждом воркер-хосте установлен компонент DuckDB CLI. Под капотом провайдер запускает DuckDB CLI, передает SQL-запрос с параметрами в CLI, получает результаты и отправляет их Airflow.
Требования
Для работы с провайдером необходимы:
-
Airflow версии 3.2.1 или выше;
-
компонент DuckDB CLI на воркер-хостах Airflow;
-
предварительно настроенное соединение Airflow для подключения к базе данных DuckDB.
Архитектура
Провайдер следует стандартной архитектуре провайдеров Airflow и включает следующие модули:
-
hooks. Хуки предоставляют интерфейсы подключения и взаимодействия с DuckDB. -
operators. Операторы выполняют запросы к DuckDB в виде задач Airflow. -
sensors. Сенсоры ожидают заданного условия в DuckDB перед запуском последующей обработки. -
utils. Модуль содержит вспомогательные функции для работы с SQL, парсинга результатов, генерации ошибок и прочие.
DuckDbHook
Провайдер содержит единственный хук — класс DuckDbHook, в котором реализована основная логика работы провайдера.
Операторы и сенсоры делегируют DuckDbHook подключение, предварительные проверки, запуск CLI, логирование и обработку ошибок.
При каждом запуске хук выполняет следующие действия:
-
Считывает объект соединения Airflow, получает путь к файлу базы данных, а также параметры поля Extra объекта соединения.
-
Проверяет
cli_paramsв списке запрещенных параметров и выполняет предварительную проверку DuckDB CLI. -
Проверяет наличие файла базы данных.
-
Передает SQL-выражение в DuckDB CLI.
-
Передает результаты запроса в Airflow.
-
Записывает stdout/stderr в логи задачи и при наличии ошибок генерирует исключения.
Основные методы DuckDbHook:
-
Отправляет SQL-строку на выполнение в DuckDB CLI. После необязательной текстовой подстановки
%(name)sхук записывает SQL во временный файл и запускает его с флагом -f, что позволяет выполнять большие SQL-запросы. Возвращает необработанный stdout без начальных и конечных пробелов в виде строки. Формат вывода по умолчанию — JSON. -
Выполняет существующий файл .sql с флагом -f. Файл выполняется без изменений и подстановки параметров. Формат вывода по умолчанию не задан, поэтому необходимо явно передать
output_format="[json|csv]". Если указанный файл .sql не найден, возвращается ошибка конфигурации. -
Проверяет доступность исполняемого файла DuckDB CLI и Extra-поля объекта соединения, включая запрещенные параметры CLI и секреты. Выполняет тестовый запрос к базе данных в памяти. Не проверяет наличие целевого файла .duckdb и возможность записи в него.
-
Помечает хук как остановленный и завершает активную группу процессов DuckDB. Используется операторами и сенсорами при остановке задачи.
DuckDbOperator
Единственный оператор провайдера DuckDbOperator выполняет SQL через DuckDB CLI на узле воркера и возвращает необработанный stdout CLI, который затем отправляется в хранилище XCom.
Основные методы DuckDbOperator:
-
Выполняет SQL-выражения и возвращает необработанную строку stdout DuckDB CLI, которую Airflow затем помещает в хранилище XCom. Метод не возвращает строки таблицы в разобранном виде.
При значении
output_format="json"по умолчанию рекомендуется использоватьairflow.providers.arenadata.duckdb.utils.json_output.parse_json_output— данный метод обрабатывает возможные префиксы, не относящиеся к JSON, и ожидает, что результат представляет собой список строк в формате JSON. Например:from airflow.providers.arenadata.duckdb.utils.json_output import parse_json_output raw = context["ti"].xcom_pull(task_ids="select_rows") rows = parse_json_output(raw)Используйте одно SQL-выражение на задачу, поскольку
_salvage_json()сохраняет лишь первый JSON-массив и отбрасывает последующие данные с предупреждением в логе. Обычныйjson.loads()безопасно использовать только для валидного JSON, представленного в виде списка. -
Завершает активный процесс DuckDB, если метод
execute()запустил хук.
DuckDbSqlSensor
Сенсор DuckDbSqlSensor ожидает, пока запрос DuckDB не вернет истинное значение.
SQL выполняется экземпляром DuckDbHook.
Сенсор анализирует первую ячейку первой возвращенной строки и продолжает ждать, если:
-
результирующий набор пуст, включая пустой stdout CLI;
-
первая строка отсутствует, не является сопоставлением (not a mapping) или не содержит столбцов;
-
первая ячейка является ложной с точки зрения Python (
false,0,nullили пустая строка).
|
ПРИМЕЧАНИЕ
Непустые строки, например "0", считаются истинными.
|
Основной метод сенсора — poke().
При каждом вызове он выполняет SQL-выражение с output_format="json", парсит JSON-результат, проверяет первую ячейку первой строки на истинность по правилам Python bool(value) и возвращает False (продолжает ожидание) или True (условие выполнено).
По умолчанию ошибки CLI и JSON приводят к сбою задачи: ненулевые коды завершения процесса DuckDB CLI и некорректный JSON вызывают исключения. Это поведение можно изменить с помощью специальных флагов класса BaseSensorOperator:
-
silent_fail=True— исключение вpoke()не приводит к сбою задачи. Оно записывается в лог, и сенсор ожидает следующий вызовpoke(); -
never_fail=True— исключение вpoke()пропускает задачу; -
soft_fail=True— ошибки CLI/JSON не влияют на ожидание или пропуск.
|
ПРИМЕЧАНИЕ
Запрос к отсутствующей таблице (например, Catalog Error: Table … does not exist) при стандартных значениях флагов немедленно завершает задачу с ошибкой, а сенсор прекращает ожидание.
Чтобы сенсор продолжал ожидать появления данных, заранее создайте таблицу (можно пустую) и используйте условие для проверки, например SELECT COUNT(*) > 0 ….
|
Сенсор парсит только первый JSON-массив в stdout, а последующие данные отбрасывает с предупреждением в логе. Поэтому SQL-выражения перед проверяемым условием не должны возвращать строки. Стоит учитывать следующее:
-
Выражения
INSTALL,LOAD,ATTACH,SETничего не выводят и могут предшествовать условию. -
CREATE SECRETвозвращает строку об успешном выполнении, которая будет воспринята как первый результирующий набор. Сенсор завершится успешно, не проверив фактическое условие, а в логе появится предупреждение. -
Каждая задача выполняется в отдельном процессе CLI. Операции
LOAD/ATTACHдолжны присутствовать в SQL сенсора, аINSTALL,CREATE PERSISTENT SECRETиCREATE VIEWможно однократно выполнить в предыдущихDuckDbOperator.
Если fail_on_empty=True и запрос не возвращает строки, сенсор вызывает AirflowFailException.
При soft_fail=True это приводит к пропуску проверки.
Если fail_on_empty=False (значение по умолчанию), при пустом результате сенсор продолжает ожидание.
Агрегатные запросы, например SELECT count(*) …, всегда возвращают одну строку, поэтому их результат не считается пустым.
Для таких запросов проверяйте истинность первой ячейки и не полагайтесь на fail_on_empty.
Каждый вызов poke() создает хук с lock_retry_attempts=0.
Ожидание разблокировки файла .duckdb или появления данных реализуется через poke_interval и параметр mode="reschedule", а не через механизм перезапуска внутри poke().
Настройка подключения
Для подключения к базе данных DuckDB провайдеру требуется предварительно настроенный объект соединения Airflow. Создайте его в веб-интерфейсе Airflow и укажите следующие параметры.
| Поле | Описание |
|---|---|
Connection type |
Тип соединения, используемый провайдером. Укажите |
Database file path |
Абсолютный путь к файлу .duckdb на узле воркера или |
Дополнительные поля
| Поле | Описание |
|---|---|
DuckDB binary |
Путь к исполняемому файлу DuckDB CLI, используемому для выполнения SQL-запросов. Путь по умолчанию (/usr/bin/duckdb) указывает на исполняемый файл, установленный компонентом DuckDB CLI сервиса DuckDB |
Дополнительные параметры JSON
Дополнительные параметры выполнения можно указать в виде JSON в поле Extra.
| Параметр | Описание | Значение по умолчанию |
|---|---|---|
timeout |
Тайм-аут завершения подпроцесса DuckDB CLI в секундах |
300 |
readonly |
Устанавливает соединение с базой данных в режиме чтения. Выполняет базовые проверки доступности файла для чтения |
false |
cli_params |
Дополнительные параметры DuckDB CLI.
Значением должна быть shell-строка ( |
— |
lock_retry_attempts |
Количество повторных попыток доступа к файлу базы данных, если он заблокирован |
0 |
{
"timeout": 300,
"readonly": false,
"cli_params": "",
"lock_retry_attempts": 0
}
Блокировка файлов и повторные попытки
При открытии файла базы данных DuckDB устанавливает эксклюзивную блокировку файла (exclusive lock).
Если файл уже заблокирован другим процессом, CLI завершается с ошибкой, в которой отражены duckdb_conn_id и db_path.
Свойство lock_retry_attempts
Для более гибкого управления блокировками используется параметр подключения lock_retry_attempts, заданный в поле Extra.
Его значение должно быть положительным целым числом, которое определяет количество повторных попыток в случае блокировки файла базы данных.
При установке значения 0 задача немедленно завершается с ошибкой и выводит сообщение о блокировке.
Интервал между повторными попытками увеличивается по схеме 1, 2, 4, 8, 16 секунд (максимальный интервал составляет 16 секунд).
Общее время выполнения задачи
Для каждой попытки применяется отдельный тайм-аут.
При включенных повторных попытках общее время выполнения задачи может достигать , поэтому параметр задачи execution_timeout следует указывать с учетом данных задержек.
Блокировки при ATTACH
Механизм повторных попыток при блокировке предполагает, что CLI устанавливает соединение с базой данных, расположенной на одном хосте.
SQL-скрипты, содержащие операции ATTACH 'file.db', выполняемые в середине процесса, могут столкнуться с теми же маркерами блокировки уже после выполнения некоторых SQL-операторов.
Для SQL-скриптов с командами ATTACH используйте lock_retry_attempts=0 (значение по умолчанию).
Логирование
В следующей таблице описаны уровни логирования для типичных событий.
| Логируемая информация | Уровень логирования |
|---|---|
Полная команда CLI |
DEBUG |
Общие параметры запуска ( |
INFO |
stdout |
INFO, обрезается до |
stderr при успешном выполнении |
INFO, обрезается до |
stderr при ошибке |
ERROR |
Повторная попытка после конфликта блокировки |
WARNING |
Запрещенные параметры CLI
Большинство встроенных параметров DuckDB, таких как выполнение SQL (-c, -f, -s), форматы вывода (-json, -csv) и другие, обрабатываются провайдером.
Указание этих параметров в cli_params приводит к ошибке конфигурации.
-
-c -
-s -
-f -
-cmd -
-init -
-json -
-csv -
-readonly -
-bail -
-no-stdin
|
РЕКОМЕНДАЦИЯ
Для указания формата вывода используйте DuckDbOperator(output_format=…).
|
Маскирование конфиденциальных данных
Учетные и другие конфиденциальные данные рекомендуется хранить в поле Extra, используя специальный набор ключей, а не непосредственно в SQL или параметрах оператора/хука.
При ошибках DuckDB обычно выводит SQL-выражение в stderr, поэтому конфиденциальные данные могут попасть в лог задачи в открытом виде.
Значения в cli_params с перечисленными ниже ключами маскируются в логах в виде ***.
-
password -
passwd -
secret -
token -
access_key -
access_token -
api_key -
apikey -
private_key -
credential -
credentials -
secret_key -
aws_secret_access_key
Например, значение токена будет скрыто в логах:
{
"cli_params": ["--token", "my-secret-value"]
}
SQL-файлы
Имена SQL-файлов должны оканчиваться на .sql (нижний регистр).
Airflow загружает SQL-скрипты из пакета DAG (dag.folder и/или template_searchpath), чтобы они были доступны и процессору DAG, и воркеру.
Используйте расширение .sql в нижнем регистре и храните ваши скрипты внутри пакета.
Если указанный SQL-файл отсутствует, возникают следующие ошибки:
Failed to resolve template field 'sql' jinja2.TemplateNotFound
Особенности жизненного цикла
DuckDbHook запускает CLI с параметром start_new_session=True, чтобы при остановке задачи или превышении тайм-аута можно было завершить группу процессов, включая процесс ADO wrapper и дочерний процесс DuckDB.
Если же воркер завершен командой SIGKILL (например, OOM или команда docker kill), метод on_kill() не вызывается, и брошенный (orphaned) процесс DuckDB может удерживать блокировку файла, пока не будет освобожден системой.
Примеры
Провайдер содержит несколько примеров DAG, демонстрирующих основные операции. Для их запуска создайте соединение Airflow в соответствии со значениями в DAG (Connection ID, Host).
Основные операции
Следующий DAG демонстрирует базовые операции провайдера DuckDB: подключение к базе данных, создание таблицы, запись и чтение тестовых данных.
from __future__ import annotations
from datetime import datetime, timedelta
from airflow.providers.arenadata.duckdb.operators.duckdb import DuckDbOperator
from airflow.providers.arenadata.duckdb.version_compat import DAG
CONN_ID = "duckdb_default" (1)
default_args = {
"owner": "airflow",
"retries": 1,
"retry_delay": timedelta(minutes=1),
}
with DAG(
dag_id="example_duckdb_basic",
start_date=datetime(2024, 1, 1),
default_args=default_args,
schedule=None,
catchup=False,
tags=["example", "duckdb"],
) as dag:
create_table = DuckDbOperator( (2)
task_id="create_table",
sql="CREATE OR REPLACE TABLE demo(id INT, name VARCHAR);",
duckdb_conn_id=CONN_ID,
)
insert_data = DuckDbOperator( (3)
task_id="insert_data",
sql="INSERT INTO demo VALUES (1, 'alpha'), (2, 'beta');",
duckdb_conn_id=CONN_ID,
)
select_count = DuckDbOperator( (4)
task_id="select_count",
sql="SELECT count(*) AS c FROM demo",
duckdb_conn_id=CONN_ID,
)
create_table >> insert_data >> select_count (5)
| 1 | Идентификатор соединения Airflow для доступа к базе данных. |
| 2 | Задача создания таблицы. |
| 3 | Задача записи данных в таблицу. |
| 4 | Задача чтения строк из таблицы. |
| 5 | Цепочка зависимостей задач. |
Работа с DuckDbSqlSensor
from __future__ import annotations
from datetime import datetime, timedelta
from airflow.providers.arenadata.duckdb.sensors.duckdb import DuckDbSqlSensor
from airflow.providers.arenadata.duckdb.version_compat import DAG, Param
CONN_ID = "duckdb_default" (1)
default_args = {
"owner": "airflow",
"retries": 1,
"retry_delay": timedelta(minutes=1),
}
with DAG(
dag_id="example_duckdb_sensors1",
start_date=datetime(2024, 1, 1),
default_args=default_args,
schedule=None,
catchup=False,
tags=["example", "duckdb", "sensor"],
params={
"table": Param("events", type="string"),
"min_id": Param(1, type="integer"),
},
) as dag:
wait_inline = DuckDbSqlSensor( (2)
task_id="wait_inline",
sql="SELECT count(*) AS ready FROM {{ params.table }}",
duckdb_conn_id=CONN_ID,
mode="reschedule",
poke_interval=30, (3)
timeout=300,
)
wait_from_sql_file = DuckDbSqlSensor(
task_id="wait_from_sql_file",
sql="queries/wait_until_ready.sql", (4)
duckdb_conn_id=CONN_ID,
mode="reschedule",
poke_interval=30,
timeout=300,
)
wait_inline >> wait_from_sql_file
| 1 | Идентификатор подключения Airflow к базе данных DuckDB. |
| 2 | Задача, ожидающая, пока SQL-выражение не вернет истинный результат (не False, 0, null, или пустую строку). |
| 3 | Интервал между проверками. |
| 4 | Задача, ожидающая, пока SQL из файла не вернет истинный результат. |
Переопределение пути к базе данных
Путь к файлу базы данных .duckdb указывается при создании соединения Airflow. Однако этот путь можно переопределить при создании экземпляра оператора. Например:
select_count = DuckDbOperator(
task_id="select_count",
sql="SELECT count(*) AS c FROM demo",
database="/tmp/example_duckdb_sql_file.duckdb", (1)
duckdb_conn_id="duckdb_default",
output_format="json",
)
| 1 | Переопределяет путь к базе данных, указанный в соединении Airflow. |