Провайдер 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, логирование и обработку ошибок.

При каждом запуске хук выполняет следующие действия:

  1. Считывает объект соединения Airflow, получает путь к файлу базы данных, а также параметры поля Extra объекта соединения.

  2. Проверяет cli_params в списке запрещенных параметров и выполняет предварительную проверку DuckDB CLI.

  3. Проверяет наличие файла базы данных.

  4. Передает SQL-выражение в DuckDB CLI.

  5. Передает результаты запроса в Airflow.

  6. Записывает stdout/stderr в логи задачи и при наличии ошибок генерирует исключения.

Основные методы DuckDbHook:

  • run_cli().

    Отправляет SQL-строку на выполнение в DuckDB CLI. После необязательной текстовой подстановки %(name)s хук записывает SQL во временный файл и запускает его с флагом -f, что позволяет выполнять большие SQL-запросы. Возвращает необработанный stdout без начальных и конечных пробелов в виде строки. Формат вывода по умолчанию — JSON.

  • run_file().

    Выполняет существующий файл .sql с флагом -f. Файл выполняется без изменений и подстановки параметров. Формат вывода по умолчанию не задан, поэтому необходимо явно передать output_format="[json|csv]". Если указанный файл .sql не найден, возвращается ошибка конфигурации.

  • test_connection()

    Проверяет доступность исполняемого файла DuckDB CLI и Extra-поля объекта соединения, включая запрещенные параметры CLI и секреты. Выполняет тестовый запрос к базе данных в памяти. Не проверяет наличие целевого файла .duckdb и возможность записи в него.

  • on_kill()

    Помечает хук как остановленный и завершает активную группу процессов DuckDB. Используется операторами и сенсорами при остановке задачи.

DuckDbOperator

Единственный оператор провайдера DuckDbOperator выполняет SQL через DuckDB CLI на узле воркера и возвращает необработанный stdout CLI, который затем отправляется в хранилище XCom.

Основные методы DuckDbOperator:

  • execute()

    Выполняет 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, представленного в виде списка.

  • on_kill()

    Завершает активный процесс 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

Тип соединения, используемый провайдером. Укажите DuckDB

Database file path

Абсолютный путь к файлу .duckdb на узле воркера или :memory: для таблиц в памяти

Дополнительные поля

Поле Описание

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-строка ("--threads 4") или JSON-массив строк (["--threads", "4"]). Нестроковые элементы в массиве вызывают ошибку конфигурации

 — 

lock_retry_attempts

Количество повторных попыток доступа к файлу базы данных, если он заблокирован

0

Пример JSON для поля Extra
{
  "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

Общие параметры запуска (duckdb_conn_id, db_path, код возврата, длительность)

INFO

stdout

INFO, обрезается до log_output_limit (именованный аргумент хука/оператора/сенсора, по умолчанию 2000). Полный текст выводится на уровне DEBUG

stderr при успешном выполнении

INFO, обрезается до log_output_limit (именованный аргумент хука/оператора/сенсора, по умолчанию 2000). Полный текст выводится на уровне DEBUG

stderr при ошибке

ERROR

Повторная попытка после конфликта блокировки

WARNING

Запрещенные параметры CLI

Большинство встроенных параметров DuckDB, таких как выполнение SQL (-c, -f, -s), форматы вывода (-json, -csv) и другие, обрабатываются провайдером. Указание этих параметров в cli_params приводит к ошибке конфигурации.

Запрещенные флаги CLI
  • -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: подключение к базе данных, создание таблицы, запись и чтение тестовых данных.

example_duckdb_basic.py
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

example_duckdb_sensors.py
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.
Нашли ошибку? Выделите текст и нажмите Ctrl+Enter чтобы сообщить о ней