You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.30 16:00:54