Snowflake 中的数据管道自动化
Emily Melhuish
Technical Curriculum Developer, Snowflake
导出到云存储
COPY INTO 直接写入 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;
| 格式 | 最适合 |
|---|---|
| CSV | 通用——几乎所有系统可读;理想的合作方导出格式 |
| JSON | 保留嵌套结构;适合半结构化消费者 |
| Parquet | 大型分析数据集;列式压缩=文件更小、读取更快 |
关键卸载选项
-- 将大导出拆分为多个文件(字节)
HEADER = TRUE -- 首行包含列名
OVERWRITE = TRUE -- 替换该路径下的现有文件
MAX_FILE_SIZE = 104857600 -- 每个输出文件 100 MB

Kafka 连接器
# connector.properties (Kafka 设置)
snowflake.topic2table.map=events:
delivery_events
snowflake.ingestion.method=
SNOWPIPE_STREAMING
Spark 连接器
// 从 Snowflake 读取为 Spark DataFrame
val df = spark.read.format("snowflake")
.options(sfOptions).option("dbtable",
"shipments").load()
通用连接
| 集成 | 类型 | 在 Harbr 的用法 |
|---|---|---|
| JDBC / ODBC | 通用驱动 | BI 工具(Tableau、Power BI、Looker)—直接查询 |
| Python Connector | 原生 Python 驱动 | 数据管道、计划 ETL、数据科学流程 |
| dbt | SQL 转换 | 在 Snowflake 计算中直接运行模型 |
| Fivetran / Airbyte | 托管摄取 | SaaS 源 → Snowflake,无需自定义代码 |
Snowflake 中的数据管道自动化