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

如何在Celery任务中运行异步函数?FastAPI场景下序列化报错求助

问题分析与解决

错误根源

你遇到的kombu.exceptions.EncodeError本质是Celery任务参数无法被JSON序列化:

  1. 传递给Celery任务的uploadchunck_func是类实例方法,方法对象不属于JSON可序列化类型;
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 04:20:23