Tự động hóa Data Pipeline trong Snowflake
Emily Melhuish
Technical Curriculum Developer, Snowflake
Xuất dữ liệu ra lưu trữ đám mây
COPY INTO ghi kết quả trực tiếp vào stageCOPY INTO @harbr_partner_export/
daily_summary/
FROM (SELECT * FROM
logistics.shipment_summary
WHERE export_date =
CURRENT_DATE() - 1);

COPY INTO @harbr_partner_export/shipment_summary/
FROM (
SELECT shipment_id,
origin,
destination,
delivery_status,
delivery_time_hours
FROM logistics.shipments
WHERE delivery_date = CURRENT_DATE()
)
FILE_FORMAT = (TYPE = 'CSV' HEADER = TRUE)
OVERWRITE = TRUE;
| Định dạng | Phù hợp nhất |
|---|---|
| CSV | Phổ biến — hầu hết hệ thống đọc được; lý tưởng để xuất cho đối tác |
| JSON | Giữ cấu trúc lồng nhau; hữu ích cho bên tiêu thụ bán cấu trúc |
| Parquet | Tập dữ liệu phân tích lớn; nén dạng cột = file nhỏ hơn, đọc nhanh hơn |
Tùy chọn unload chính
-- Chia export lớn thành nhiều file (bytes)
HEADER = TRUE -- gồm tên cột ở dòng đầu
OVERWRITE = TRUE -- thay thế file đã có tại đường dẫn
MAX_FILE_SIZE = 104857600 -- 100 MB mỗi file

Kafka Connector
# connector.properties (Kafka settings)
snowflake.topic2table.map=events:
delivery_events
snowflake.ingestion.method=
SNOWPIPE_STREAMING
Spark Connector
// Đọc từ Snowflake vào Spark DataFrame
val df = spark.read.format("snowflake")
.options(sfOptions).option("dbtable",
"shipments").load()
Kết nối phổ quát
| Tích hợp | Loại | Dùng tại Harbr |
|---|---|---|
| JDBC / ODBC | Trình điều khiển phổ quát | Công cụ BI (Tableau, Power BI, Looker) - truy vấn trực tiếp |
| Python Connector | Trình điều khiển Python gốc | Pipeline dữ liệu, ETL theo lịch, quy trình khoa học dữ liệu |
| dbt | Biến đổi SQL | Chạy model trực tiếp trên compute của Snowflake |
| Fivetran / Airbyte | Ingestion quản lý | Nguồn SaaS → Snowflake, không cần code tùy chỉnh |
Tự động hóa Data Pipeline trong Snowflake