如何在Celery任务中运行异步函数?FastAPI场景下序列化报错求助
问题分析与解决
错误根源
你遇到的kombu.exceptions.EncodeError本质是Celery任务参数无法被JSON序列化:
- 传递给Celery任务的
uploadchunck_func是类实例方法,方法对象不属于JSON可序列化类型; - FastAPI的
UploadFile是绑定当前请求上下文的对象,同样无法被序列化传递到Celery Worker进程。
async_to_sync不是问题所在,即使转换了异步函数,参数无法序列化的问题依然存在。
解决方案
重构代码,避免直接传递不可序列化的对象,改为传递可序列化的数据(如文件内容、路径、字符串/UUID等),并在Celery任务内部直接调用处理逻辑。
方案1:传递文件二进制内容(适合小文件)
修改Celery任务,直接在任务内部实例化服务类/调用处理函数:
from asgiref.sync import async_to_sync from your_module import UploadChunkService # 导入你的上传处理服务 @app.task def upload_zip_file(dataset_id: UUID, file_content: bytes, filename: str, user_id: UUID): upload_service = UploadChunkService() # 调用异步方法并转换为同步执行 duplicate_count, duplicate_filenames = async_to_sync(upload_service.group_send)( dataset_id, file_content, filename, user_id ) return duplicate_count, duplicate_filenames
在FastAPI接口中读取文件内容后调用任务:
async def update_dataset_upload2(dataset_id: UUID, file: UploadFile, user_id: UUID): # 读取文件二进制内容 file_content = await file.read() # 调用Celery任务,传递可序列化参数 task = upload_zip_file.delay( dataset_id=dataset_id, file_content=file_content, filename=file.filename, user_id=user_id ) return {"task_id": task.id}
方案2:传递文件路径(适合大文件)
如果文件体积较大,不适合直接传递二进制内容,先将文件保存到临时存储,再传递路径给Celery任务:
FastAPI接口代码:
import tempfile from pathlib import Path async def update_dataset_upload2(dataset_id: UUID, file: UploadFile, user_id: UUID): # 创建临时文件保存上传内容 with tempfile.NamedTemporaryFile(delete=False, suffix=Path(file.filename).suffix) as tmp_file: content = await file.read() tmp_file.write(content) tmp_path = tmp_file.name # 调用Celery任务 task = upload_zip_file.delay( dataset_id=dataset_id, file_path=tmp_path, filename=file.filename, user_id=user_id ) return {"task_id": task.id}
Celery任务代码:
from asgiref.sync import async_to_sync from your_module import UploadChunkService from pathlib import Path @app.task def upload_zip_file(dataset_id: UUID, file_path: str, filename: str, user_id: UUID): upload_service = UploadChunkService() # 读取临时文件内容 with open(file_path, 'rb') as f: file_content = f.read() duplicate_count, duplicate_filenames = async_to_sync(upload_service.group_send)( dataset_id, file_content, filename, user_id ) # 清理临时文件 Path(file_path).unlink() return duplicate_count, duplicate_filenames
关键注意事项
- Celery任务的所有参数必须是JSON可序列化类型(如str、int、bytes、UUID、dict等),禁止传递类实例、方法、请求上下文对象;
- 如果
group_send方法依赖UploadFile对象的结构,需要修改该方法,使其能够接受二进制内容+文件名的组合,而不是直接依赖UploadFile。
内容的提问来源于stack exchange,提问作者Quân nguyễn
相关产品推荐
相关产品推荐

