Luồng và Change Data Capture

Tự động hóa Data Pipeline trong Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Chỉ xử lý phần thay đổi

CDC: Change Data Capture Hai cách so sánh cách dữ liệu được xử lý

Tự động hóa Data Pipeline trong Snowflake

Streams

Stream làm gì

  • Theo dõi mọi INSERT, UPDATE, DELETE trên bảng nguồn

Screenshot 2026-05-11 at 10.50.31 am.png

  • Duy trì nhật ký thay đổi liên tục — không sao chép dữ liệu
  • Khi đã tiêu thụ, offset tiến lên; lần đọc sau bắt đầu mới
1 * Tài liệu học Snowflake
Tự động hóa Data Pipeline trong Snowflake

Các loại Stream

 

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)

Tự động hóa Data Pipeline trong Snowflake

Tạo Stream

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;
Tự động hóa Data Pipeline trong Snowflake

Các cột siêu dữ liệu của Stream

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTION: INSERT hoặc DELETE
  • METADATA$ISUPDATE: TRUE khi là một phần của cặp cập nhật
  • METADATA$ROW_ID: Mã định danh hàng vật lý duy nhất
  • Cập nhật xuất hiện thành cặp DELETE + INSERT, cả hai đều có cờ METADATA$ISUPDATE = TRUE

Các cột siêu dữ liệu của stream hiển thị METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID với dữ liệu mẫu

Tự động hóa Data Pipeline trong Snowflake

Offset của Stream

Sơ đồ dòng thời gian - Tạo Stream (offset bắt đầu tại đây) → Nguồn có thay đổi (stream tích lũy bản ghi) → Stream được tiêu thụ trong một giao dịch (offset tiến đến thời điểm hiện tại)

Tự động hóa Data Pipeline trong Snowflake

Tổng quan: Stream trong Pipeline

  • Stream ghép với task — đối tượng Snowflake chạy SQL theo lịch
  • Task chỉ đọc hàng đã thay đổi; với 10M hàng: 2 giây so với 2 phút

Screenshot 2026-05-11 at 10.48.51 am.png

1 * Tài liệu học Snowflake
Tự động hóa Data Pipeline trong Snowflake

Truy vấn: Stream trong Pipeline

  • Stream ghép với task - đối tượng Snowflake chạy SQL theo lịch
  • Task chỉ đọc hàng đã thay đổi; với 10M hàng: 2 giây so với 2 phút
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

Hãy thực hành!

Tự động hóa Data Pipeline trong Snowflake

Preparing Video For Download...