如何在AirFlow中用SqlSensor基于更新时间获取增量记录
解决AirFlow SqlSensor时间戳变量及更新问题
针对你的需求,这里提供两种可行方案,分别适配不同的业务场景:
方案1:基于DAG调度周期的时间戳(固定间隔批量处理)
适合你的@hourly固定调度场景,利用AirFlow内置的模板变量直接获取调度时间节点:
代码修改
- 替换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持久化上次处理的最新时间戳,每次轮询时从该时间点开始查询:
代码修改
- SQL中读取XCom存储的时间:首次运行时 fallback 到DAG的
start_date - 成功函数更新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
相关产品推荐
相关产品推荐

