Snowpipe 与 Snowpipe Streaming

Snowflake 中的数据管道自动化

Emily Melhuish

Technical Curriculum Developer, Snowflake

Snowpipe

用例:

  • 交付事件持续到达

解决方案:Snowpipe

  • Snowpipe 在文件出现在 stage 后立即从文件加载数据。

截图 2026-05-11 下午12.24.14

1 * Snowflake 学习材料
Snowflake 中的数据管道自动化

批量加载的问题

  • 文件全天每隔几分钟到达 S3
  • 计划的 COPY INTO 午夜运行——滞后 24 小时
  • 延迟的发货要到次日才会出现
  • Snowpipe 弥补这一缺口
-- 每夜批处理:00:00 运行,数据全天到达
COPY INTO logistics.delivery_events
FROM @harbr_s3_stage/events/
FILE_FORMAT = (FORMAT_NAME = 'harbr_json_format');
-- 9 点的异常要到明天才会出现
Snowflake 中的数据管道自动化

什么是 Snowpipe?

  • 封装一条 COPY INTO 语句——语法与文件格式相同
  • 当新文件到达 stage 时自动触发
  • 以微批方式加载,通常在几分钟内完成
  • 无服务器——无需预配仓库
CREATE PIPE harbr_events_pipe AS
  COPY INTO logistics.delivery_events
  FROM @harbr_s3_stage/events/
  FILE_FORMAT = (FORMAT_NAME = 'harbr_json_format');
Snowflake 中的数据管道自动化

Snowpipe 的工作原理

Snowpipe 工作流

  • AUTO_INGEST——事件驱动;云存储发布通知
  • Amazon S3 | Azure Event Grid | GCP Pub/Sub
  • REST API 触发——从编排代码直接调用 insertFilesinsertReport 端点
Snowflake 中的数据管道自动化

Snowpipe 计费

截图 2026-05-11 下午12.24.14

  • 按每消耗 1 GB 的固定信用额度计费
  • 文本文件:按未压缩大小计费
  • 二进制文件:按观测大小计费
Snowflake 中的数据管道自动化

Snowpipe Streaming

Snowpipe Snowpipe Streaming
触发 文件到达 stage 应用写入行
延迟 分钟级 秒级
用途 基于文件的事件流 GPS、IoT、实时应用数据

 

完全消除文件边界

  • 通过 Streaming Ingest SDK 从应用直接写入行
  • 无文件、无 stage——延迟为秒级
# Snowpipe Streaming:应用直接写入行
channel = client.openChannel('GPS_CHANNEL', 'LOGISTICS', 'GPS_EVENTS')
channel.insertRows(rows=[
    {'vehicle_id': 'V001', 'lat': 51.5, 'lng': -0.12, 'ts': now()}
])
Snowflake 中的数据管道自动化

选择合适的摄取方法

摄取方法

方法 何时使用
COPY INTO 定时批量加载:每日文件、每周导出;可接受小时级延迟
Snowpipe 文件持续到达;需在到达后数分钟内加载
Snowpipe Streaming 应用生成数据:GPS、IoT、金融市场;数据秒级可用
1 * Snowflake 学习资源
Snowflake 中的数据管道自动化

让我们练习吧!

Snowflake 中的数据管道自动化

Preparing Video For Download...