स्ट्रीम्स और चेंज डेटा कैप्चर

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

Emily Melhuish

Technical Curriculum Developer, Snowflake

केवल बदला हुआ डेटा प्रोसेस करना

CDC: Change Data Capture डेटा प्रोसेसिंग के दो तुलना

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

स्ट्रीम्स

स्ट्रीम क्या करती है

  • सोर्स टेबल पर हर INSERT, UPDATE, DELETE ट्रैक करती है

Screenshot 2026-05-11 at 10.50.31 am.png

  • रनिंग चेंज लॉग बनाए रखती है — डेटा डुप्लीकेशन नहीं
  • कंज्यूम होने पर ऑफसेट आगे बढ़ता है; अगली रीड नई से शुरू होती है
1 * Snowflake Learning Resource
Snowflake में डेटा पाइपलाइन ऑटोमेशन

स्ट्रीम प्रकार

 

स्ट्रीम प्रकार क्या कैप्चर करता है उपयुक्त उपयोग
स्टैंडर्ड सभी टेबल प्रकार और व्यूज़; सभी DML बदलाव — इन्सर्ट, अपडेट, डिलीट ट्रैक करता है जहाँ कोई भी रो बदल सकती है (जैसे shipments)
एपेंड-ओनली सभी टेबल प्रकार और व्यूज़, बाहरी टेबल छोड़कर — केवल रो इन्सर्ट ट्रैक एक-बार-इन्सर्ट टेबल (जैसे delivery events) — अधिक कुशल
इन्सर्ट-ओनली बाहरी रूप से प्रबंधित Apache Iceberg और एक्सटर्नल टेबल — केवल रो इन्सर्ट ट्रैक एक्सटर्नल टेबल

डायरेक्टरी टेबल्स किसी स्टेज के फ़ाइल मेटाडेटा दिखाती हैं (नाम, आकार, अंतिम संशोधन टाइमस्टैम्प)

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

स्ट्रीम बनाना

shipments टेबल पर स्टैंडर्ड स्ट्रीम

CREATE STREAM shipments_stream
  ON TABLE logistics.shipments;

delivery events टेबल पर एपेंड-ओनली स्ट्रीम

CREATE STREAM delivery_events_stream
  ON TABLE logistics.delivery_events
  APPEND_ONLY = TRUE;
Snowflake में डेटा पाइपलाइन ऑटोमेशन

स्ट्रीम मेटाडेटा कॉलम

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTION: INSERT या DELETE
  • METADATA$ISUPDATE: अपडेट जोड़ी का भाग होने पर TRUE
  • METADATA$ROW_ID: यूनिक फिजिकल रो पहचानकर्ता
  • अपडेट्स एक DELETE + INSERT जोड़ी के रूप में दिखते हैं; दोनों में METADATA$ISUPDATE = TRUE होता है

स्ट्रीम मेटाडेटा कॉलम: METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID के साथ नमूना डेटा

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

स्ट्रीम ऑफसेट

टाइमलाइन आरेख - स्ट्रीम बनाई गई (ऑफसेट यहाँ से शुरू) → सोर्स टेबल में बदलाव होते हैं (स्ट्रीम रिकॉर्ड जोड़ती है) → ट्रांज़ैक्शन में स्ट्रीम उपभोग की गई (ऑफसेट अब तक बढ़ता है)

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

पाइपलाइन में स्ट्रीम्स: अवलोकन

  • स्ट्रीम्स टास्क्स के साथ पेयर होती हैं — Snowflake ऑब्जेक्ट्स जो शेड्यूल पर SQL चलाते हैं
  • टास्क केवल बदली हुई पंक्तियाँ पढ़ता है; 10M पंक्तियों पर: 2 सेकंड बनाम 2 मिनट

Screenshot 2026-05-11 at 10.48.51 am.png

1 * Snowflake Learning Resource
Snowflake में डेटा पाइपलाइन ऑटोमेशन

पाइपलाइन में स्ट्रीम्स: क्वेरी

  • स्ट्रीम्स टास्क्स के साथ पेयर होती हैं - Snowflake ऑब्जेक्ट्स जो शेड्यूल पर SQL चलाते हैं
  • टास्क केवल बदली हुई पंक्तियाँ पढ़ता है; 10M पंक्तियों पर: 2 सेकंड बनाम 2 मिनट
CREATE TASK logistics.sync_shipments
  WAREHOUSE = compute_wh
  SCHEDULE = '5 MINUTE'
  WHEN SYSTEM$STREAM_HAS_DATA('logistics.staging_shipments_stream')
AS
  INSERT INTO logistics.shipments
  SELECT shipment_id, region, carrier, delivery_days
  FROM logistics.staging_shipments_stream
  WHERE METADATA$ACTION = 'INSERT';
Snowflake में डेटा पाइपलाइन ऑटोमेशन

Ayo berlatih!

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

Preparing Video For Download...