如何在Airflow中串行调度DAG以批量处理数据库作业?
问题
数据库中生成了一批作业,我有一个可处理单个作业的计算密集型DAG,但资源十分有限。如何调度DAG使它们串行运行以处理所有作业?
我最初的方案是每2分钟运行一个调度DAG,流程为:[检查无作业正在处理 -> 触发处理DAG -> 标记作业已处理],但感觉不符合Airflow的惯用写法。这类任务通常用任务队列服务(如Amazon SQS、GCP Cloud Tasks)解决,但如果只能使用Airflow该怎么做?
优化后的解决方案
更新调度DAG后,Airflow可保证同一时间仅运行一个该调度DAG实例,方案更简洁且自定义代码更少:
流程为:[取作业 -> 触发处理DAG -> 标记作业已处理]
代码示例
with DAG( max_active_runs=1, schedule_interval='* * * * *' ) as dag: # 编写具体任务逻辑:take_job、trigger_processing_dag、mark_job_as_processed
参数说明
max_active_runs=1:指定Airflow同一时间仅运行一个该调度DAG的实例schedule_interval='* * * * *':若当前无运行实例则周期性触发该DAG
内容的提问来源于stack exchange,提问作者vimi
相关产品推荐
相关产品推荐

