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

如何在AirFlow中用SqlSensor基于更新时间获取增量记录

解决AirFlow SqlSensor时间戳变量及更新问题

针对你的需求,这里提供两种可行方案,分别适配不同的业务场景:


方案1:基于DAG调度周期的时间戳(固定间隔批量处理)

适合你的@hourly固定调度场景,利用AirFlow内置的模板变量直接获取调度时间节点:

代码修改

  1. 替换SQL中的时间变量:用{{ prev_execution_date }}(上一次DAG执行时间)作为查询起始时间,第一次运行时会默认使用DAG的start_date:
with DAG(
        dag_id="dag_process_supervisor",
        start_date=datetime(2022, 12, 1),
        catchup=False,
        schedule_interval="@hourly"
) as dag:
    # 定义成功判断函数:检查是否有返回记录
    def _success_criteria(record):
        return len(record) > 0

    wait_for_table_update = SqlSensor(
        task_id='forecasting_jobs_sensor',
        conn_id='postgres',
        sql='''
          SELECT *
          FROM ForecastingJob 
          WHERE last_modified_at > '{{ prev_execution_date }}';
       ''',
        success=_success_criteria,
        pass_value=True,
        timeout=5 * 60,
        poke_interval=60,
        mode='reschedule'
    )

适用场景

每小时批量处理上一个调度周期内新增/更新的记录,适合对实时性要求不高的批量业务。


方案2:用XCom存储上次处理时间(实时轮询跟踪)

适合需要实时检测新增记录的场景,通过XCom持久化上次处理的最新时间戳,每次轮询时从该时间点开始查询:

代码修改

  1. SQL中读取XCom存储的时间:首次运行时 fallback 到DAG的start_date
  2. 成功函数更新XCom时间:检测到记录后,将最新的last_modified_at存入XCom供下次使用
with DAG(
        dag_id="dag_process_supervisor",
        start_date=datetime(2022, 12, 1),
        catchup=False,
        schedule_interval="@hourly"
) as dag:
    # 定义带上下文的成功判断函数
    def _success_criteria(record, ti):
        if len(record) == 0:
            return False
        # 获取本次查询到的最新更新时间
        max_last_modified = max(row['last_modified_at'] for row in record)
        # 将时间存入XCom,供下次轮询使用
        ti.xcom_push(key='last_processed_time', value=max_last_modified.isoformat())
        return True

    wait_for_table_update = SqlSensor(
        task_id='forecasting_jobs_sensor',
        conn_id='postgres',
        sql='''
          SELECT *
          FROM ForecastingJob 
          WHERE last_modified_at > COALESCE('{{ ti.xcom_pull(task_ids="forecasting_jobs_sensor", key="last_processed_time") }}', '{{ dag.start_date }}');
       ''',
        success=_success_criteria,
        pass_value=True,
        timeout=5 * 60,
        poke_interval=60,
        mode='reschedule',
        provide_context=True  # 必须开启,让函数能获取TaskInstance对象
    )

注意事项

  • 确保last_modified_at字段为带时区的timestamptz类型,避免时区偏移问题
  • 定期清理AirFlow元数据库中的旧XCom数据,避免存储冗余

补充说明

SqlSensor仅负责等待条件满足,如果你需要拉取并处理这些记录,需在传感器后添加后续任务(比如PostgresOperator或PythonOperator)来执行数据处理逻辑。

内容的提问来源于stack exchange,提问作者Mohamed Maherz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 18:45:30