Stream dan Change Data Capture

Otomatisasi Data Pipeline di Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Memproses Hanya yang Berubah

CDC: Change Data Capture Dua perbandingan tentang bagaimana data diproses

Otomatisasi Data Pipeline di Snowflake

Stream

Fungsi stream

  • Melacak setiap INSERT, UPDATE, dan DELETE pada tabel sumber

Tangkapan layar 2026-05-11 pukul 10.50.31

  • Menjaga log perubahan berjalan — tanpa duplikasi data
  • Setelah dikonsumsi, offset maju; pembacaan berikutnya mulai baru
1 * Sumber Belajar Snowflake
Otomatisasi Data Pipeline di Snowflake

Jenis Stream

 

Jenis Stream Menangkap Terbaik Untuk
Standar Semua jenis tabel dan view & semua perubahan DML - melacak insert, update, delete Tabel di mana baris mana pun bisa berubah (mis. pengiriman)
Append-only Semua jenis tabel dan view, kecuali tabel eksternal - hanya melacak insert baris Tabel sekali-insert (mis. peristiwa pengantaran) - lebih efisien
Insert-only Apache Iceberg dan tabel eksternal yang dikelola eksternal - hanya melacak insert baris Tabel eksternal

Tabel direktori menampilkan metadata file untuk stage (nama, ukuran, cap waktu modifikasi terakhir)

Otomatisasi Data Pipeline di Snowflake

Membuat Stream

Stream standar pada tabel shipments

CREATE STREAM shipments_stream
  ON TABLE logistics.shipments;

Stream append-only pada tabel delivery events

CREATE STREAM delivery_events_stream
  ON TABLE logistics.delivery_events
  APPEND_ONLY = TRUE;
Otomatisasi Data Pipeline di Snowflake

Kolom Metadata Stream

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTION: INSERT atau DELETE
  • METADATA$ISUPDATE: TRUE saat bagian dari pasangan pembaruan
  • METADATA$ROW_ID: Pengidentifikasi baris fisik unik
  • Update muncul sebagai pasangan DELETE + INSERT, keduanya ditandai METADATA$ISUPDATE = TRUE

Kolom metadata stream menampilkan kolom METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID dengan data contoh

Otomatisasi Data Pipeline di Snowflake

Offset Stream

Diagram lini masa - Stream dibuat (offset mulai di sini) → Perubahan terjadi di tabel sumber (stream mengakumulasi rekaman) → Stream dikonsumsi dalam transaksi (offset maju ke sekarang

Otomatisasi Data Pipeline di Snowflake

Ikhtisar Stream dalam Pipeline

  • Stream dipasangkan dengan task — objek Snowflake yang menjalankan SQL terjadwal
  • Task hanya membaca baris yang berubah; dengan 10M baris: 2 detik vs 2 menit

Screenshot 2026-05-11 at 10.48.51 am.png

1 * Sumber Belajar Snowflake
Otomatisasi Data Pipeline di Snowflake

Kueri Stream dalam Pipeline

  • Stream dipasangkan dengan task - objek Snowflake yang menjalankan SQL terjadwal
  • Task hanya membaca baris yang berubah; dengan 10M baris: 2 detik vs 2 menit
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';
Otomatisasi Data Pipeline di Snowflake

Ayo berlatih!

Otomatisasi Data Pipeline di Snowflake

Preparing Video For Download...