FastAPI 0.105.0集成Celery5.3.6触发ValueError错误求助
问题:CentOS7下FastAPI集成Celery触发ValueError: too many values to unpack (expected 2)
在CentOS 7环境中使用FastAPI 0.105.0集成Celery 5.3.6时,调用Celery任务(无论是apply_async还是delay)都会触发ValueError: too many values to unpack (expected 2)错误,错误触发点在任务调用行。
错误栈
/root/anaconda3/envs/fastapi_rebuild/bin/python /root/Desktop/fastapi_rebuild/main.py /root/Desktop/fastapi_rebuild/main.py:22: DeprecationWarning: on_event is deprecated, use lifespan event handlers instead. Read more about it in the [FastAPI docs for Lifespan Events](https://fastapi.tiangolo.com/advanced/events/). @app.on_event("startup") INFO: Started server process [25185] INFO: Waiting for application startup. INFO: Application startup complete. INFO: Uvicorn running on http://0.0.0.0:9001 (Press CTRL+C to quit) INFO: 192.168.18.2:62283 - "POST /mask_file HTTP/1.1" 500 Internal Server Error ERROR: Exception in ASGI application Traceback (most recent call last): File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/uvicorn/protocols/http/h11_impl.py", line 408, in run_asgi result = await app( # type: ignore[func-returns-value] File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/uvicorn/middleware/proxy_headers.py", line 84, in __call__ return await self.app(scope, receive, send) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/fastapi/applications.py", line 1106, in __call__ await super().__call__(scope, receive, send) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/starlette/applications.py", line 122, in __call__ await self.middleware_stack(scope, receive, send) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/starlette/middleware/errors.py", line 184, in __call__ raise exc File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/starlette/middleware/errors.py", line 162, in __call__ await self.app(scope, receive, _send) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/starlette/middleware/exceptions.py", line 79, in __call__ raise exc File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/starlette/middleware/exceptions.py", line 68, in __call__ await self.app(scope, receive, sender) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/fastapi/middleware/asyncexitstack.py", line 20, in __call__ raise e File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/fastapi/middleware/asyncexitstack.py", line 17, in __call__ await self.app(scope, receive, send) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/starlette/routing.py", line 718, in __call__ await route.handle(scope, receive, send) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/starlette/routing.py", line 276, in handle await self.app(scope, receive, send) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/starlette/routing.py", line 66, in app response = await func(request) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/fastapi/routing.py", line 274, in app raw_response = await run_endpoint_function( File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/fastapi/routing.py", line 191, in run_endpoint_function return await dependant.call(**values) File "/root/Desktop/fastapi_rebuild/router/mask_file_router.py", line 54, in test_run_task task_id = docx_task.apply_async((file.filename, keywords, upload_method, mask_method, key, minio_path)) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/celery/app/task.py", line 594, in apply_async return app.send_task( File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/celery/app/base.py", line 728, in send_task router = router or amqp.router File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/kombu/utils/objects.py", line 40, in __get__ return super().__get__(instance, owner) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/functools.py", line 993, in __get__ val = self.func(instance) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/celery/app/amqp.py", line 581, in router return self.Router() File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/celery/app/amqp.py", line 264, in Router return _routes.Router(self.routes, queues or self.queues, File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/celery/app/amqp.py", line 576, in routes self.flush_routes() File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/celery/app/amqp.py", line 269, in flush_routes self._rtable = _routes.prepare(self.app.conf.task_routes) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/celery/app/routes.py", line 136, in prepare return [expand_route(route) for route in routes] File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/celery/app/routes.py", line 136, in <listcomp> return [expand_route(route) for route in routes] File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/celery/app/routes.py", line 127, in expand_route return MapRoute(route) File "/root/anaconda3/envs/fastapi_rebuild/lib/python3.9/site-packages/celery/app/routes.py", line 33, in __init__ for k, v in map: ValueError: too many values to unpack (expected 2)
相关代码与配置
Celery配置(celery_app.py)
celery = Celery(__name__) celery.conf.broker_url = os.getenv("BROKER_URL") celery.conf.result_backend = os.getenv("RESULT_BACKEND") celery.conf.task_routes = ([ ("worker_queue.gevent_tasks.*", {"queue": "gevent_queue"}), ("worker_queue.prefork_tasks.*", {"queue": "prefork_queue"}) ])
Worker实现(workers/docx_file.py)
def mask_docx_file(upload_file, mask_words, upload_method, mask_method, key, minio_path): return "111"
Celery任务(worker_queue/gevent_tasks.py)
@celery.task def docx_task(upload_file, mask_words, upload_method, mask_method, key, minio_path): return mask_docx_file(upload_file, mask_words, upload_method, mask_method, key, minio_path)
任务调用代码
if file_type == 'docx': task_id = docx_task.apply_async((file.filename, keywords, upload_method, mask_method, key, minio_path)) return task_id.id
项目目录结构
fastapi_rebuild/ └── celery_task/ ├── celery_utils/ │ ├── __init__.py │ ├── encoding_format.py │ ├── key_processor.py │ └── minio_operation.py ├── worker_queue/ │ ├── __init__.py │ ├── gevent_tasks.py │ └── prefork_tasks.py ├── workers/ │ ├── __init__.py │ └── docx_file.py ├── celery_app.py
原因分析
从错误栈可以追踪到,问题出在Celery解析task_routes配置时:
- Celery 5.x版本中,
task_routes如果使用列表格式,列表中的每个元素必须是包含单个键值对的字典; - 当前配置用了列表包裹元组的形式,Celery在解析时试图将元组拆成
k, v,但元组的第二个元素本身是字典,导致拆包时出现"too many values to unpack"错误。
解决方案
修改Celery的task_routes配置为以下两种正确格式之一:
方式1:列表包裹字典
celery.conf.task_routes = [ {"worker_queue.gevent_tasks.*": {"queue": "gevent_queue"}}, {"worker_queue.prefork_tasks.*": {"queue": "prefork_queue"}} ]
方式2:直接使用字典格式(推荐)
celery.conf.task_routes = { "worker_queue.gevent_tasks.*": {"queue": "gevent_queue"}, "worker_queue.prefork_tasks.*": {"queue": "prefork_queue"} }
修改后重启FastAPI服务和Celery Worker,即可解决该错误。
内容的提问来源于stack exchange,提问作者刘佳诚
相关产品推荐
相关产品推荐

