Задачи и оркестрация DAG

Автоматизация конвейеров данных в Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Решение для ручных рабочих процессов: задачи

  • Команда вручную запускает SQL-скрипт каждое утро
  • Один пропуск — бизнес теряет видимость
  • Задачи Snowflake полностью автоматизируют процесс

Пример:

-- Today's manual routine (fragile)
CALL logistics.refresh_ops_dashboard();
Автоматизация конвейеров данных в Snowflake

Основные сведения о задачах

  • Задачи выполняют SQL по расписанию — внешний CRON не нужен
  • Выполняют SQL-операторы или вызывают хранимые процедуры
  • Вычисления и расписание встроены в Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Автоматизация конвейеров данных в Snowflake

Создание отдельной задачи

CREATE OR REPLACE TASK transform_delivery_summary 
  WAREHOUSE = harbr_wh 
  SCHEDULE = 'USING CRON 5 * * * * UTC' 
  AS INSERT INTO delivery_summary 

SELECT shipment_id, status, updated_at 
FROM delivery_events 
WHERE processed = FALSE; 
Автоматизация конвейеров данных в Snowflake

Склад или бессерверный режим

На основе склада

  • Использует именованный виртуальный склад
  • Минимум 60 секунд тарификации за запуск
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

Бессерверная

  • Пропустите условие WAREHOUSE
  • Оплата посекундно; простой не тарифицируется
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Автоматизация конвейеров данных в Snowflake

Оркестрация задач на основе DAG

Схема DAG — ingest_raw (корневая задача, значок часов) → clean_events (выполняется после ingest_raw) → build_summary (выполняется после clean_events)

  • DAG: направленный ациклический граф
  • Корневая задача содержит расписание CRON
  • Дочерние задачи указывают предшественника через AFTER
  • Вся цепочка централизована
  • DAG управляет порядком и временем выполнения
Автоматизация конвейеров данных в Snowflake

Управление состояниями задач

Диаграмма состояний — SUSPENDED → STARTED → SUCCEEDED или FAILED

-- Activate a task
ALTER TASK mytask RESUME; 
-- Pause a task
ALTER TASK mytask SUSPEND; 
-- Execute a task
EXECUTE TASK mytask SUSPEND;
Автоматизация конвейеров данных в Snowflake

Совместное использование задач и потоков

Схема CDC-конвейера — delivery_events → поток фиксирует изменения → есть данные? ДА → задача запускается → delivery_summary обновляется / НЕТ → задача пропускается

CREATE OR REPLACE TASK logistics.process_delivery_events
    WAREHOUSE = compute_wh
    SCHEDULE = '5 minute'
    WHEN SYSTEM$STREAM_HAS_DATA('logistics.delivery_events_stream')
AS
INSERT INTO logistics.processed_events
SELECT * FROM logistics.delivery_events_stream;
Автоматизация конвейеров данных в Snowflake

Давайте попрактикуемся!

Автоматизация конвейеров данных в Snowflake

Preparing Video For Download...