使用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"]
排查与解决方案
核心问题定位
chain任务未启动:process_ml_model_dataset_factor_x仅构建了chain对象,但未调用执行(如result_chain()),导致process_ml_model_chunk任务仅被发送到Broker,实际执行链路未触发。- 主任务阻塞占用资源:
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
相关产品推荐
相关产品推荐

