कार्य और DAG ऑर्केस्ट्रेशन

Snowflake में डेटा पाइपलाइन ऑटोमेशन

Emily Melhuish

Technical Curriculum Developer, Snowflake

मैनुअल वर्कलोड का समाधान: Tasks

  • टीम हर सुबह SQL स्क्रिप्ट हाथ से चलाती है
  • एक रन छूटा — बिज़नेस की विज़िबिलिटी घटती है
  • Snowflake tasks इसे पूरी तरह स्वचालित करते हैं

उदाहरण:

-- आज की मैनुअल दिनचर्या (नाज़ुक)
CALL logistics.refresh_ops_dashboard();
Snowflake में डेटा पाइपलाइन ऑटोमेशन

Tasks के बारे में क्या जानें

  • टास्क शेड्यूल पर 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 में डेटा पाइपलाइन ऑटोमेशन

Standalone टास्क बनाना

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 क्लॉज़ छोड़ें
  • उपभोगित प्रति सेकंड बिल; idle लागत नहीं
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 के बाद चलता है)

  • DAGs: Directed Acyclic Graph
  • रूट टास्क में CRON शेड्यूल होता है
  • चाइल्ड टास्क AFTER से अपना पूर्ववर्ती बताते हैं
  • पूरी चेन केंद्रीकृत है
  • DAGs तय करते हैं कब और किस क्रम में चीजें चलें
Snowflake में डेटा पाइपलाइन ऑटोमेशन

टास्क स्टेट प्रबंधन

स्टेट आरेख — SUSPENDED → STARTED → SUCCEEDED या FAILED

-- किसी टास्क को सक्रिय करें
ALTER TASK mytask RESUME; 
-- किसी टास्क को रोकें
ALTER TASK mytask SUSPEND; 
-- किसी टास्क को निष्पादित करें
EXECUTE TASK mytask SUSPEND;
Snowflake में डेटा पाइपलाइन ऑटोमेशन

टास्क और स्ट्रीम एक साथ

CDC पाइपलाइन आरेख — delivery_events → स्ट्रीम बदलाव कैप्चर करती है → स्ट्रीम में डेटा? YES → टास्क ट्रिगर → delivery_summary अपडेट / NO → टास्क स्किप

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 में डेटा पाइपलाइन ऑटोमेशन

Ayo berlatih!

Snowflake में डेटा पाइपलाइन ऑटोमेशन

Preparing Video For Download...