任务与 DAG 编排

Snowflake 中的数据管道自动化

Emily Melhuish

Technical Curriculum Developer, Snowflake

手工工作负载的解决方案:任务

  • 团队每天早上手动运行一段 SQL 脚本
  • 漏跑一次,业务可视性就受损
  • 使用 Snowflake 任务可完全自动化

示例:

-- 当天的手工流程(脆弱)
CALL logistics.refresh_ops_dashboard();
Snowflake 中的数据管道自动化

任务要点

  • 任务按计划运行 SQL——无需外部 CRON
  • 可执行 SQL 语句或调用存储过程
  • 计算与调度原生运行在 Snowflake 中
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Snowflake 中的数据管道自动化

创建独立任务

CREATE OR REPLACE TASK transform_delivery_summary 
  WAREHOUSE = harbr_wh 
  SCHEDULE = 'USING CRON 5 * * * * UTC' 
  AS INSERT INTO delivery_summary 

SELECT shipment_id, status, updated_at 
FROM delivery_events 
WHERE processed = FALSE; 
Snowflake 中的数据管道自动化

基于仓库 vs 无服务器

基于仓库(Warehouse)

  • 使用命名虚拟仓库
  • 每次运行最少计费 60 秒
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

无服务器(Serverless)

  • 省略 WAREHOUSE 子句
  • 按实际秒计费;无空闲成本
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Snowflake 中的数据管道自动化

基于 DAG 的任务编排

DAG 示意图 — ingest_raw(根任务,时钟图标)→ clean_events(在 ingest_raw 之后运行)→ build_summary(在 clean_events 之后运行)

  • DAG:有向无环图
  • 根任务承载 CRON 调度
  • 子任务用 AFTER 声明其前置任务
  • 整条链集中管理
  • DAG 控制运行时间与顺序
Snowflake 中的数据管道自动化

管理任务状态

状态图 — SUSPENDED → STARTED → SUCCEEDED 或 FAILED

-- 启用任务
ALTER TASK mytask RESUME; 
-- 暂停任务
ALTER TASK mytask SUSPEND; 
-- 执行任务
EXECUTE TASK mytask SUSPEND;
Snowflake 中的数据管道自动化

任务与流的组合

CDC 管道示意 — delivery_events → 流捕获变更 → 流有数据?是 → 触发任务 → 更新 delivery_summary / 否 → 跳过任务

CREATE OR REPLACE TASK logistics.process_delivery_events
    WAREHOUSE = compute_wh
    SCHEDULE = '5 minute'
    WHEN SYSTEM$STREAM_HAS_DATA('logistics.delivery_events_stream')
AS
INSERT INTO logistics.processed_events
SELECT * FROM logistics.delivery_events_stream;
Snowflake 中的数据管道自动化

Ayo berlatih!

Snowflake 中的数据管道自动化

Preparing Video For Download...