如何通过Python API在Snowflake DAG中检查Stream是否有数据
解决Snowflake DAG中通过Python API检查Stream数据的问题
核心方案
在通过Python API创建Snowflake DAG任务时,只需将WHEN SYSTEM$STREAM_HAS_DATA()条件直接整合到目标任务的SQL定义中,就能实现仅当Stream存在数据时才执行任务的逻辑。
具体实现
方案一:直接修改任务SQL定义
把条件语句添加到INSERT任务的SQL开头,确保任务执行前先检查Stream数据:
CONNECTION_PARAMETERS = {"account": "XXX", "user": "YYY", "password": "", "warehouse":"TEST_WH", "database": "TEST_DB", "schema": "TEST_SCHEMA"} connection = connect(**CONNECTION_PARAMETERS) root = Root(connection) dag = DAG("TEMP_DAG", schedule=timedelta(minutes=1), warehouse="TEST_WH") with dag: test_parent_task = DAGTask( "TEST_PARENT_TASK", definition=""" WHEN SYSTEM$STREAM_HAS_DATA('TEST_DB.TEST_SCHEMA.TEST_STREAM') INSERT INTO TEST_INTERMEDIATE(ID,NAME,LASTUPDATED_TIMESTAMP,INGESTION_TIME,TOPIC) SELECT ID ,NAME, current_timestamp(), RECORD_INGESTION_TIME, TOPIC FROM TEST_DB.TEST_SCHEMA.TEST_STREAM; """, warehouse="TEST_WH" ) test_child_task = DAGTask( "TEST_CHILD_TASK", definition="call DUMMY_PROC()", warehouse="TEST_WH" ) test_child_task >> test_parent_task schema = root.databases["TEST_DB"].schemas["TEST_SCHEMA"] dag_op = DAGOperation(schema) dag_op.deploy(dag)
方案二:动态拼接SQL(适用于多环境场景)
如果需要根据传入的数据库、架构参数动态生成检查条件,可以先拼接SQL字符串再传入DAGTask:
def test(database,schema_env,warehouse,test_config,root): dag = DAG("TESTDAG", schedule=timedelta(minutes=1), warehouse=warehouse) # 拼接完整的Stream标识符 stream_full_path = f"{database}.{schema_env}.TEST_STREAM" # 组合带检查条件的任务SQL parent_task_def = f""" WHEN SYSTEM$STREAM_HAS_DATA('{stream_full_path}') {test_config["PROD_PARENT_TASK"]} """ with dag: STREAM_TO_INTERMEDIATE = DAGTask( test_config["PARENT_TASK_NAME"], parent_task_def, warehouse=warehouse ) INTERMEDIATE_TO_STAGE = DAGTask( test_config["CHILD_TASK_NAME"], test_config["CHILD_TASK"], warehouse=warehouse ) INTERMEDIATE_TO_STAGE >> STREAM_TO_INTERMEDIATE schema = root.databases[database].schemas[schema_env] dag_op = DAGOperation(schema) dag_op.deploy(dag)
关键注意点
SYSTEM$STREAM_HAS_DATA()必须传入完整的Stream路径(数据库.架构.流名称),避免因当前会话上下文导致的识别错误。- 该条件是Snowflake任务的原生控制逻辑,通过Python API部署后,Snowflake会自动在任务执行前完成数据检查,无需额外编写Python层面的校验代码。
内容的提问来源于stack exchange,提问作者Liya
相关产品推荐
相关产品推荐

