Otomatisasi Data Pipeline di Snowflake
Emily Melhuish
Technical Curriculum Developer, Snowflake
CDC: Change Data Capture

Fungsi 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)
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;
SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
METADATA$ACTION: INSERT atau DELETEMETADATA$ISUPDATE: TRUE saat bagian dari pasangan pembaruanMETADATA$ROW_ID: Pengidentifikasi baris fisik unikMETADATA$ISUPDATE = TRUE


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