Tự động hóa Data Pipeline trong Snowflake
Emily Melhuish
Technical Curriculum Developer, Snowflake
CDC: Change Data Capture

Stream làm gì

| Loại stream | Ghi nhận | Phù hợp nhất |
|---|---|---|
| Standard | Mọi loại bảng và view & mọi thay đổi DML - theo dõi insert, update, delete | Bảng có thể thay đổi bất kỳ hàng nào (vd. shipments) |
| Append-only | Mọi loại bảng và view, trừ external tables - chỉ theo dõi insert | Bảng chỉ chèn (vd. delivery events) - hiệu quả hơn |
| Insert-only | Apache Iceberg do bên ngoài quản lý và external tables - chỉ theo dõi insert | External tables |
Directory tables hiển thị siêu dữ liệu tệp cho một stage (tên, kích thước, thời điểm sửa đổi cuối)
Standard stream trên bảng shipments
CREATE STREAM shipments_stream
ON TABLE logistics.shipments;
Append-only stream trên bảng 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 hoặc DELETEMETADATA$ISUPDATE: TRUE khi là một phần của cặp cập nhậtMETADATA$ROW_ID: Mã định danh hàng vật lý duy nhấtMETADATA$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';
Tự động hóa Data Pipeline trong Snowflake