Airflow 2.5.0使用billiard多进程任务卡住,求解决方案
Airflow 2.5.0中billiard多进程卡住问题解决及正确实现方式
问题原因
你遇到的卡住问题核心原因有两个:
- 若使用CeleryExecutor,Airflow本身就是基于billiard实现worker进程的,在任务内部嵌套创建billiard Pool会引发进程通信冲突,导致子进程无法正常接收任务。
- 任务内部定义的
test是嵌套函数,无法被billiard的序列化机制正确传递给子进程,进而导致进程挂起无响应。
现有代码的修复方案
针对你提供的代码,做以下修改即可解决卡住问题:
- 把嵌套的
test函数移到任务函数外部,确保能被正常序列化 - 显式管理Pool的生命周期(避免用上下文管理器,减少潜在的进程回收问题)
修改后的代码:
import pendulum from billiard import Pool from airflow import DAG from airflow.decorators import task # 将函数移到外部,保证可序列化 def test(l): return sum(l) with DAG(dag_id='ttest', schedule_interval="40 * * * *", start_date=pendulum.datetime(2023, 1, 1, tz="UTC")) as dag: @task(task_id='te') def test_task(ds=None, **kwargs): a = [[1,2], [2,3], [3,4]] print('start pool') pool = Pool(2) try: res = pool.map(test, a) print(res) finally: # 显式关闭并等待进程结束 pool.close() pool.join() t = test_task()
Airflow任务中多进程的正确实现方式
Airflow的核心设计是任务级并行,而非单个任务内的进程并行,优先拆分任务更符合Airflow的调度逻辑,以下是几种推荐方式:
1. 拆分为独立子任务
把需要并行处理的单元拆成单独的Airflow Task,利用Airflow本身的Executor(如Celery、Kubernetes)实现并行:
import pendulum from airflow import DAG from airflow.decorators import task def calculate_sum(l): return sum(l) with DAG(dag_id='parallel_tasks', schedule_interval="40 * * * *", start_date=pendulum.datetime(2023, 1, 1, tz="UTC")) as dag: input_data = [[1,2], [2,3], [3,4]] # 动态生成并行子任务 for idx, data in enumerate(input_data): @task(task_id=f'calc_sum_{idx}') def sum_task(data): result = calculate_sum(data) print(result) return result sum_task(data)
2. 用TaskGroup管理批量并行任务
如果需要统一管理一组并行任务,使用TaskGroup会让DAG结构更清晰:
import pendulum from airflow import DAG from airflow.decorators import task, task_group def calculate_sum(l): return sum(l) with DAG(dag_id='parallel_task_group', schedule_interval="40 * * * *", start_date=pendulum.datetime(2023, 1, 1, tz="UTC")) as dag: @task_group(group_id='sum_calculations') def sum_task_group(): input_data = [[1,2], [2,3], [3,4]] for idx, data in enumerate(input_data): @task(task_id=f'calc_{idx}') def task_func(data): return calculate_sum(data) task_func(data) sum_task_group()
3. 特殊场景下的任务内并行(不推荐)
如果业务逻辑必须在单个任务内做并行计算,推荐使用concurrent.futures.ProcessPoolExecutor,同时确保所有传递给子进程的内容可序列化:
import pendulum from concurrent.futures import ProcessPoolExecutor from airflow import DAG from airflow.decorators import task def test(l): return sum(l) with DAG(dag_id='in_task_parallel', schedule_interval="40 * * * *", start_date=pendulum.datetime(2023, 1, 1, tz="UTC")) as dag: @task(task_id='te') def test_task(ds=None, **kwargs): a = [[1,2], [2,3], [3,4]] print('start parallel') with ProcessPoolExecutor(max_workers=2) as executor: res = list(executor.map(test, a)) print(res) t = test_task()
关键注意事项
- 优先使用Airflow的任务级并行,避免在单个任务内嵌套多进程,这会增加调试难度,且不符合Airflow的设计初衷
- 任何任务内的多进程实现,必须确保传递给子进程的函数、数据都能被序列化(禁止用嵌套函数、lambda、无法序列化的自定义对象)
- 如果使用CeleryExecutor,不要嵌套使用billiard,因为Celery本身已经用billiard管理worker进程,嵌套会引发进程通信异常
内容的提问来源于stack exchange,提问作者viminal
相关产品推荐
相关产品推荐

