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

Airflow Master Dag调度Dataflow并行子任务时SQL Server表未写入问题排查

问题根因分析与解决方案

首先明确:该问题确实和DagSensor(即Airflow的ExternalDagSensor)的配置错误直接相关,同时存在任务依赖配置的逻辑问题,结合你给出的场景,具体原因和修复方案如下:

1. 核心错误:任务依赖顺序错误

你给出的代码中任务依赖写法是错误的:

task1_pipeline >> [trigger_dag_Abcd, status_check_abcd]

该写法意味着触发子Dag的任务和DagSensor会同时启动,DagSensor会从任务启动时刻就开始轮询目标子Dag的运行状态。而Airflow的ExternalDagSensor默认会匹配和当前Master Dag执行时间一致的最近一次子Dag运行记录,如果同一调度周期之前已经有过成功的子Dag运行记录,Sensor会直接返回成功,哪怕本次刚触发的子Dag还在运行、甚至还没完成启动。
这就会导致后续的SQL Server加载任务提前执行,此时API类子Dag还没有生成最新的CSV文件,加载任务要么读取到旧的空文件、要么读取到不完整的文件,自然没有新数据写入SQL Server。而手动单独执行加载任务时,CSV已经完全生成,所以运行正常。

修复方案

调整任务依赖为串行顺序,先触发子Dag,再启动Sensor轮询:

task1_pipeline >> trigger_dag_Abcd >> status_check_abcd

2. ExternalDagSensor的额外配置坑

就算调整了依赖顺序,还要补充两个配置避免偶发异常:

  • 增加mode="reschedule"参数:默认的poke模式会一直占用worker槽位,7个Sensor并行运行时如果Composer的worker资源不足,会导致部分Sensor长时间无法启动,超时后直接标记失败/误判状态。reschedule模式会在每次轮询间隔释放worker资源,大幅降低资源占用。
  • 确认执行时间匹配规则:如果Master Dag和子Dag的调度周期不一致,需要补充execution_date_fn参数自定义匹配的子Dag执行时间,避免Sensor匹配到历史的子Dag成功记录,误判当前任务状态。

调整后的Sensor示例:

DagSensor(
        dag=dag,
        task_id='status_check_abcd',
        external_dag_id='ABCD_DF',
        poke_interval=30,
        timeout=4000,
        mode="reschedule"
)

3. 并行加载的补充排查点

如果调整完Sensor配置后仍有偶发失败,需要排查两个常见问题:

  • 检查7个加载任务是否写入同一张表,是否存在表级排他锁冲突,部分Dataflow任务遇到锁冲突时会默认静默跳过写入,仅标记任务成功
  • 检查API子Dag生成CSV后是否有明确的GCS文件落盘完成标记,避免加载任务读取到还没完全上传的临时文件

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 22:36:03