Оркестрація завдань і DAG

Автоматизація конвеєрів даних у Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Рішення для ручних навантажень: завдання

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

Приклад:

-- Сьогоднішня ручна рутина (крихка)
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

-- Активувати завдання
ALTER TASK mytask RESUME; 
-- Призупинити завдання
ALTER TASK mytask SUSPEND; 
-- Виконати завдання
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...