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

使用Python+Flask+Celery+Redis时多Group任务无法执行含Chain的子任务

问题描述

基于Nickjj的docker-flask-example项目,使用Python、Flask、Celery和Redis构建任务系统时出现以下异常:

  • 当调用包含多个子任务的group(n_factor>1)时,嵌套chain的process_ml_model_dataset_factor_x任务无法执行
  • Worker日志显示已接收process_ml_model_chunk任务,但从未实际执行
  • 当n_factor=1(仅单个任务的group)时,所有任务均正常运行

已尝试Celery 5.3.6和5.4.0rc2版本,相关代码与配置如下:

任务代码片段

@celery.task(bind=True, ignore_result=False)
def process_ml_model_dataset(self, dataset_serialized, media_serialized, ml_model_serialized, collection_name, n_factor):

    if n_factor == 1:
        result_group = group(process_ml_model_dataset_factor_x.s(dataset_serialized, media_serialized, ml_model_serialized, collection_name, n_factor))
    else:
        result_group = (group(process_ml_model_dataset_factor_x.s(dataset_serialized, media_serialized, ml_model_serialized, collection_name, i) for i in range(1, n_factor+1)))

    result = result_group()

    while not result.ready():
        time.sleep(60)

@celery.task(bind=True, ignore_result=False)
def process_ml_model_dataset_factor_x(self, dataset_serialized, media_serialized, ml_model_serialized, collection_name, factor_x):
    result_chain = chain(
        process_ml_model_chunk.s(None, 1, dataset_serialized, ml_model_serialized, first_chunk_pickle, collection_name, factor_x))

    reader_second_chunk = FILE_HELPER.csv_reader_chunked(
        MEDIAS_UPLOADS_DEFAULT_DEST + "/" + str(media_serialized["id"]) + "/" + media_serialized["file_name"],
        columns_header, DATASET_CHUNK_SIZE, DATASET_CHUNK_SIZE)
    for i, chunk in enumerate(reader_second_chunk):
        result_chain |= process_ml_model_chunk.s(i + 2, dataset_serialized, ml_model_serialized, pickle.dumps(chunk),
                                                 collection_name, factor_x)

Worker日志

worker-1  | [2024-04-03 09:42:23,143: INFO/MainProcess] Task src.tasks.dataset_full_ml_process.process_ml_model_chunk[d3c55d4c-2c1e-48e5-927f-50dce7ac6738] received
worker-1  | [2024-04-03 09:42:23,210: INFO/MainProcess] Task src.tasks.dataset_full_ml_process.process_ml_model_chunk[26af7b45-d1fd-41f6-88c3-04b2f3dd9f6a] received
worker-1  | [2024-04-03 09:42:23,273: INFO/MainProcess] Task src.tasks.dataset_full_ml_process.process_ml_model_chunk[0314e273-a66b-493a-88e6-6ae81f30cb5e] received

Celery配置

REDIS_URL = os.getenv("REDIS_URL", "redis://redis:6379/0")

# Celery.
CELERY_CONFIG = {
    "broker_url": REDIS_URL,
    "result_backend": REDIS_URL,
    "include": [
        "src.tasks.dataset_full_ml_process"
    ],
}

Docker Worker配置

worker:
<<: *default-app
command: celery -A "src.app.celery_app" worker -l "${CELERY_LOG_LEVEL:-debug}"
entrypoint: []
deploy:
  resources:
    limits:
      cpus: "${DOCKER_WORKER_CPUS:-0}"
      memory: "${DOCKER_WORKER_MEMORY:-0}"
profiles: ["worker"]
排查与解决方案

核心问题定位

  1. chain任务未启动:process_ml_model_dataset_factor_x仅构建了chain对象,但未调用执行(如result_chain()),导致process_ml_model_chunk任务仅被发送到Broker,实际执行链路未触发。
  2. 主任务阻塞占用资源:process_ml_model_dataset中用while循环轮询等待结果,当n_factor>1时,多个子任务并发会耗尽Worker进程资源,导致后续任务无法调度。

具体修复步骤

1. 启动process_ml_model_dataset_factor_x中的chain

修改任务代码,在构建完chain后调用执行:

@celery.task(bind=True, ignore_result=False)
def process_ml_model_dataset_factor_x(self, dataset_serialized, media_serialized, ml_model_serialized, collection_name, factor_x):
    result_chain = chain(
        process_ml_model_chunk.s(None, 1, dataset_serialized, ml_model_serialized, first_chunk_pickle, collection_name, factor_x))

    reader_second_chunk = FILE_HELPER.csv_reader_chunked(
        MEDIAS_UPLOADS_DEFAULT_DEST + "/" + str(media_serialized["id"]) + "/" + media_serialized["file_name"],
        columns_header, DATASET_CHUNK_SIZE, DATASET_CHUNK_SIZE)
    for i, chunk in enumerate(reader_second_chunk):
        result_chain |= process_ml_model_chunk.s(i + 2, dataset_serialized, ml_model_serialized, pickle.dumps(chunk),
                                                 collection_name, factor_x)
    # 新增:启动chain任务执行
    result_chain()

2. 优化主任务等待逻辑

替换轮询阻塞逻辑为异步回调,避免占用Worker进程:

@celery.task(bind=True, ignore_result=False)
def process_ml_model_dataset(self, dataset_serialized, media_serialized, ml_model_serialized, collection_name, n_factor):
    tasks = [
        process_ml_model_dataset_factor_x.s(dataset_serialized, media_serialized, ml_model_serialized, collection_name, i)
        for i in range(1, n_factor+1)
    ]
    result_group = group(*tasks)
    # 用异步回调替代轮询,不阻塞Worker
    result_group.apply_async(link=after_group_completion.s())

# 新增回调任务,处理group完成后的逻辑
@celery.task
def after_group_completion(results):
    # 此处编写group任务全部完成后的处理逻辑
    pass

3. 调整Worker资源配置

Docker配置中cpus和memory设为0意味着无资源限制,易引发竞争。建议设置合理限制,并指定Worker并发数:

deploy:
  resources:
    limits:
      cpus: "${DOCKER_WORKER_CPUS:-2}"
      memory: "${DOCKER_WORKER_MEMORY:-2G}"

修改Worker启动命令,指定并发数:

command: celery -A "src.app.celery_app" worker -l "${CELERY_LOG_LEVEL:-debug}" --concurrency=4

4. 验证任务参数序列化

在process_ml_model_chunk任务开头添加日志,确认参数能正确反序列化:

@celery.task(bind=True, ignore_result=False)
def process_ml_model_chunk(self, *args, **kwargs):
    self.logger.info(f"Received chunk task with args: {args}")
    # 原有任务逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 20:36:00