งานและการจัดการ DAG

การทำให้ Data Pipeline เป็นอัตโนมัติใน Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

แก้ปัญหางานที่ทำด้วยตนเองด้วย Tasks

  • ทีมงานรัน SQL script ด้วยตนเองทุกเช้า
  • พลาดครั้งเดียว ธุรกิจสูญเสียข้อมูล
  • Snowflake tasks ทำให้กระบวนการนี้เป็นอัตโนมัติ

ตัวอย่าง:

-- Today's manual routine (fragile)
CALL logistics.refresh_ops_dashboard();
การทำให้ Data Pipeline เป็นอัตโนมัติใน Snowflake

สิ่งที่ควรรู้เกี่ยวกับ Tasks

  • Tasks รัน SQL ตามตาราง — ไม่ต้องใช้ CRON ภายนอก
  • รัน SQL statement หรือเรียก stored procedure
  • คำนวณและกำหนดเวลาภายใน Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
การทำให้ Data Pipeline เป็นอัตโนมัติใน Snowflake

การสร้าง Standalone Task

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; 
การทำให้ Data Pipeline เป็นอัตโนมัติใน Snowflake

แบบใช้ Warehouse กับแบบ Serverless

แบบใช้ Warehouse

  • ใช้ virtual warehouse ที่ระบุชื่อ
  • เรียกเก็บค่าใช้จ่ายขั้นต่ำ 60 วินาทีต่อครั้ง
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

แบบ Serverless

  • ละเว้น clause WAREHOUSE
  • เรียกเก็บตามการใช้งานจริง ไม่มีค่า idle
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
การทำให้ Data Pipeline เป็นอัตโนมัติใน Snowflake

การจัดการงานด้วย DAG

ไดอะแกรม DAG — ingest_raw (งานรากฐาน, ไอคอนนาฬิกา) → clean_events (ทำงานหลัง ingest_raw) → build_summary (ทำงานหลัง clean_events)

  • DAG: Directed Acyclic Graph
  • งานรากฐานกำหนดตาราง CRON
  • งานย่อยระบุงานก่อนหน้าด้วย AFTER
  • ควบคุมทั้งสายงานจากจุดเดียว
  • DAG กำหนดเวลาและลำดับการทำงาน
การทำให้ Data Pipeline เป็นอัตโนมัติใน 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;
การทำให้ Data Pipeline เป็นอัตโนมัติใน Snowflake

การใช้งาน Tasks ร่วมกับ Streams

ไดอะแกรม CDC pipeline — delivery_events → stream บันทึกการเปลี่ยนแปลง → stream มีข้อมูล? ใช่ → งานทำงาน → อัปเดต 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;
การทำให้ Data Pipeline เป็นอัตโนมัติใน Snowflake

ฝึกกันเลย!

การทำให้ Data Pipeline เป็นอัตโนมัติใน Snowflake

Preparing Video For Download...