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

Airflow 2.5.0使用billiard多进程任务卡住,求解决方案

Airflow 2.5.0中billiard多进程卡住问题解决及正确实现方式

问题原因

你遇到的卡住问题核心原因有两个:

  1. 若使用CeleryExecutor,Airflow本身就是基于billiard实现worker进程的,在任务内部嵌套创建billiard Pool会引发进程通信冲突,导致子进程无法正常接收任务。
  2. 任务内部定义的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 14:01:17