Snowflake 中的数据管道自动化
Emily Melhuish
Technical Curriculum Developer, Snowflake
示例:
-- 当天的手工流程(脆弱)
CALL logistics.refresh_ops_dashboard();
CREATE TASK logistics.refresh_dashboard
WAREHOUSE = harbr_wh
SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
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;
基于仓库(Warehouse)
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 ...;

AFTER 声明其前置任务
-- 启用任务
ALTER TASK mytask RESUME;
-- 暂停任务
ALTER TASK mytask SUSPEND;
-- 执行任务
EXECUTE TASK mytask SUSPEND;

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 中的数据管道自动化