Điều phối Tasks và DAG

Tự động hóa Data Pipeline trong Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Giải pháp cho tác vụ thủ công: Tasks

  • Nhóm chạy thủ công một script SQL mỗi sáng
  • Bỏ lỡ một lần chạy → mất khả năng quan sát kinh doanh
  • Snowflake tasks tự động hóa hoàn toàn

Ví dụ:

-- Quy trình thủ công hôm nay (dễ lỗi)
CALL logistics.refresh_ops_dashboard();
Tự động hóa Data Pipeline trong Snowflake

Cần biết gì về Tasks

  • Tasks chạy SQL theo lịch — không cần CRON ngoài
  • Thực thi câu lệnh SQL hoặc gọi stored procedure
  • Tính toán và lập lịch chạy nguyên bản trong Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Tự động hóa Data Pipeline trong Snowflake

Tạo Task độc lập

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; 
Tự động hóa Data Pipeline trong Snowflake

Warehouse vs Serverless

Dựa trên Warehouse

  • Dùng một virtual warehouse có tên
  • Tối thiểu tính phí 60 giây mỗi lần chạy
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

Serverless

  • Bỏ WAREHOUSE
  • Tính phí theo giây sử dụng; không có chi phí nhàn rỗi
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Tự động hóa Data Pipeline trong Snowflake

Điều phối nhiệm vụ dựa trên DAG

Sơ đồ DAG — ingest_raw (nhiệm vụ gốc, biểu tượng đồng hồ) → clean_events (chạy sau ingest_raw) → build_summary (chạy sau clean_events)

  • DAG: Directed Acyclic Graph
  • Nhiệm vụ gốc chứa lịch CRON
  • Nhiệm vụ con khai báo tiền nhiệm bằng AFTER
  • Chuỗi được quản lý tập trung
  • DAG quyết định thời điểm và thứ tự chạy
Tự động hóa Data Pipeline trong Snowflake

Quản lý trạng thái Task

Sơ đồ trạng thái — SUSPENDED → STARTED → SUCCEEDED hoặc FAILED

-- Kích hoạt task
ALTER TASK mytask RESUME; 
-- Tạm dừng task
ALTER TASK mytask SUSPEND; 
-- Thực thi task
EXECUTE TASK mytask SUSPEND;
Tự động hóa Data Pipeline trong Snowflake

Kết hợp Tasks và Streams

Sơ đồ pipeline CDC — delivery_events → stream ghi nhận thay đổi → stream có dữ liệu? CÓ → task chạy → cập nhật delivery_summary / KHÔNG → task bỏ qua

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;
Tự động hóa Data Pipeline trong Snowflake

Ayo berlatih!

Tự động hóa Data Pipeline trong Snowflake

Preparing Video For Download...