Snowflake 中的数据管道自动化
Emily Melhuish
Technical Curriculum Developer, Snowflake
用例:
解决方案:Snowpipe

COPY INTO 午夜运行——滞后 24 小时-- 每夜批处理:00:00 运行,数据全天到达
COPY INTO logistics.delivery_events
FROM @harbr_s3_stage/events/
FILE_FORMAT = (FORMAT_NAME = 'harbr_json_format');
-- 9 点的异常要到明天才会出现
COPY INTO 语句——语法与文件格式相同CREATE PIPE harbr_events_pipe AS
COPY INTO logistics.delivery_events
FROM @harbr_s3_stage/events/
FILE_FORMAT = (FORMAT_NAME = 'harbr_json_format');

AUTO_INGEST——事件驱动;云存储发布通知insertFiles 或 insertReport 端点
| Snowpipe | Snowpipe Streaming | |
|---|---|---|
| 触发 | 文件到达 stage | 应用写入行 |
| 延迟 | 分钟级 | 秒级 |
| 用途 | 基于文件的事件流 | GPS、IoT、实时应用数据 |
完全消除文件边界
# Snowpipe Streaming:应用直接写入行
channel = client.openChannel('GPS_CHANNEL', 'LOGISTICS', 'GPS_EVENTS')
channel.insertRows(rows=[
{'vehicle_id': 'V001', 'lat': 51.5, 'lng': -0.12, 'ts': now()}
])

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