Snowflake多STREAM_HAS_DATA条件任务的问题咨询与报错排查
Snowflake任务相关问题及报错修复
问题1:判断触发任务的流
任务触发后,无需嵌套IF,有两种更直接的方式判断哪个流满足条件:
- 提前存储流状态变量:在任务开始时先检查每个流的状态并存入变量,后续直接判断变量即可:
DECLARE stream1_has_data BOOLEAN := SYSTEM$STREAM_HAS_DATA('db.schema.stream_1'); stream2_has_data BOOLEAN := SYSTEM$STREAM_HAS_DATA('db.schema.stream_2'); stream3_has_data BOOLEAN := SYSTEM$STREAM_HAS_DATA('db.schema.stream_3'); BEGIN IF stream1_has_data THEN -- 处理stream1逻辑 ELSIF stream2_has_data THEN -- 处理stream2逻辑 ELSIF stream3_has_data THEN -- 处理stream3逻辑 END IF; END;
- 带标识查询多流数据:通过UNION ALL将多个流的数据合并,同时标记来源流,遍历结果时直接获取触发源:
SELECT 'stream_1' AS source_stream, process_date FROM db.schema.stream_1 UNION ALL SELECT 'stream_2' AS source_stream, process_date FROM db.schema.stream_2 UNION ALL SELECT 'stream_3' AS source_stream, process_date FROM db.schema.stream_3;
问题2:查询流是否会清空数据
不会。Snowflake流的数据只有在被消费时才会被清空,消费操作包括:
- 执行包含流的DML语句(INSERT/UPDATE/DELETE/MERGE)
- 使用
CREATE TABLE ... AS SELECT基于流创建表 - 使用
INSERT INTO ... SELECT将流数据插入目标表
单纯的SELECT查询(包括COUNT统计)不属于消费操作,不会清空流中的记录。
报错脚本修复
错误原因
- IFF函数误用:IFF是返回值的表达式函数,不能在分支中执行变量赋值操作,原代码中
IFF(..., stream_result := ..., ...)违反语法规则。 - 游标声明时机不当:初始DECLARE块中声明游标时,
stream_result尚未赋值,可能导致遍历异常。 - 未定义循环名称:原代码使用
BREAK inner_loop;但未给内部循环命名,语法不合法。 - 变量名错误:RETURN语句中引用的
:respuesta与定义的变量response_string不匹配。 - 存储过程调用语法错误:直接用
response_string := (call sp_1(...))无法正确获取存储过程返回值。
修复后的脚本
CREATE OR REPLACE TASK db.schema.task_name WAREHOUSE = XXXX USER_TASK_TIMEOUT_MS = XXXX SCHEDULE = 'USING CRON 0,30 * * * * America/Santiago' WHEN SYSTEM$STREAM_HAS_DATA('db.schema.stream_1') OR SYSTEM$STREAM_HAS_DATA('db.schema.stream_2') OR SYSTEM$STREAM_HAS_DATA('db.schema.stream_3') OR SYSTEM$STREAM_HAS_DATA('db.schema.stream_4') AS EXECUTE IMMEDIATE $$ DECLARE response_string STRING; result BOOLEAN; stream_result resultset; query_result resultset; BEGIN -- 改用IF-ELSIF实现流判断与赋值 IF SYSTEM$STREAM_HAS_DATA('db.schema.stream_1') THEN stream_result := (SELECT process_date FROM db.schema.stream_1 ORDER BY process_date); ELSIF SYSTEM$STREAM_HAS_DATA('db.schema.stream_2') THEN stream_result := (SELECT process_date FROM db.schema.stream_2 ORDER BY process_date); ELSIF SYSTEM$STREAM_HAS_DATA('db.schema.stream_3') THEN stream_result := (SELECT process_date FROM db.schema.stream_3 ORDER BY process_date); ELSIF SYSTEM$STREAM_HAS_DATA('db.schema.stream_4') THEN stream_result := (SELECT process_date FROM db.schema.stream_4 ORDER BY process_date); END IF; -- 赋值后声明游标,确保关联有效结果集 DECLARE c1 CURSOR FOR stream_result; FOR row_variable IN c1 DO query_result := (SELECT stg.flg_registry FROM db.schema.load_status stg WHERE row_variable.process_date = stg.process_date); DECLARE c2 CURSOR FOR query_result; -- 初始化结果为true,避免未赋值问题 result := TRUE; FOR row_variable_2 IN c2 DO IF row_variable_2.flg_registry != 1 THEN result := false; BREAK; END IF; END FOR; IF result THEN -- 正确调用存储过程并获取返回值 CALL sp_1(row_variable.process_date) INTO response_string; RETURN response_string; ELSE RETURN 'process date ' || row_variable.process_date || ' is not full loaded'; END IF; END FOR; END; $$;
内容的提问来源于stack exchange,提问作者Robertino Bonora
相关产品推荐
相关产品推荐

