ARTICLE DETAIL

资讯详情

深耕商务建站与企业官网运营的一线实战洞察。

从零实现带收件箱的AI助手:FastAPI异步任务实战

从零实现带收件箱的AI助手:FastAPI异步任务实战 先说说最近看到的一个有意思的项目。有人在 Hacker News 上展示了一个 AI 助手卖点不是聊天对话多顺畅而是它自带一个独立收件箱inbox。用户可以往里面投递任务助手异步消费、处理、回填结果整个流程像一套轻量级的消息队列。这种设计在 AI Agent 工程实践里越来越常见当助手不再只是“问答机器人”而需要处理批量写作、定时巡检、工单分类、内容审核等异步任务时同步聊天的模式就不够用了。本文会从零实现一个“带收件箱的 AI 助手”后端服务只依赖 FastAPI 和 Python 标准库。我们会逐步拆解收件箱的任务模型、状态流转、Worker 消费逻辑以及完整的 API 接口最后给出常见问题排查思路和生产环境落地的建议。无论你是刚开始接触 AI 工程化还是想给已有助手系统增加异步任务能力这篇教程都值得收藏。1. 背景与核心概念1.1 什么是“带收件箱的 AI 助手”传统的 AI 助手通常是一个同步聊天接口用户发送问题模型返回回答调用结束后整个交互就结束了。这个模式对“聊天”场景没问题但对任务型场景存在明显短板。比如用户一次性提交 20 封邮件让助手生成摘要如果同步处理客户端必须长时间等待一旦网络波动或接口超时前面所有结果都丢失了。“带收件箱的 AI 助手”借鉴了异步消息系统的设计思路。助手内部维护一个收件箱所有请求先进入收件箱排队后台 Worker 不断从收件箱中取出任务调用模型或工具处理再把结果回写到对应的任务记录上。用户提交请求后拿到一个task_id之后可以通过这个 ID 查询处理进度和最终结果。从架构角度看收件箱本质上是一个任务队列只是它需要额外支持任务状态管理排队中、处理中、已完成、失败。按优先级或时间排序取数。任务结果回写与查询。失败重试与错误信息记录。这套设计并不新鲜消息队列领域已经实践了多年。但当它被应用到 AI 助手场景时有一个很大的区别AI 模型调用往往是慢操作而且可能失败必须把“任务状态”和“处理结果”作为一等公民来管理。1.2 收件箱模式解决的核心问题给 AI 助手引入收件箱模式主要解决四个问题。第一是解耦。调用方只需要把任务投递到收件箱不需要关心 AI 模型什么时候处理完。后端可以随时增加消费 Worker也可以平滑升级模型服务调用方是无感的。第二是可靠性。任务被持久化后即使 Worker 进程崩溃任务记录不会丢失。重启后可以继续消费未完成的任务这比同步调用里的“请求丢失”场景可靠得多。第三是可观测性。所有任务都有明确状态我们可以方便地统计队列积压量、平均处理时长、失败率甚至对每个任务做审计。第四是并发可控。我们可以限制同时处理的任务数量避免大批量请求瞬间压垮模型 API也可以结合令牌桶做限流。下面用一个对比表来总结同步聊天与收件箱模式的区别能力维度同步聊天模式收件箱模式请求方式请求-响应提交任务-异步回调/轮询任务状态无明确状态PENDING/PROCESSING/DONE/FAILED持久性依赖客户端连接任务记录持久化并发控制较难Worker 数量可控失败重试需要客户端重试服务端自动重试适用场景实时对话批量处理、后台任务、Agent 任务编排1.3 典型应用场景在实际项目中这种模式很适合以下场景。内容生成与摘要批量生成商品文案、新闻摘要、邮件回复草稿。工单分类与回复客服工单进入收件箱AI 自动打标、分配、生成建议回复。数据处理任务从数据库或文件中抽取数据交给模型结构化再写回存储。定时巡检报告每天定时把运营数据丢进收件箱模型生成日报后推送通知。这些场景的共同点是任务到达时间和处理时间不一定是同步的而且单次处理可能耗时几十秒甚至几分钟。用收件箱模式能最大程度降低系统耦合度。2. 系统架构与消息状态设计2.1 整体架构我们设计的系统包含四个核心角色API 层、Inbox 存储层、Worker 消费层、AI 处理服务。下面用一张 ASCII 架构图表示数据流客户端 (提交任务/查询结果) │ ▼ ┌─────────────────┐ │ FastAPI 层 │ /inbox/tasks 提交 │ │ /inbox/tasks 查询 └─────────────────┘ │ ▼ ┌─────────────────┐ │ Inbox 存储 │ 任务状态 内存队列 └─────────────────┘ │ │ 领取待处理任务 ▼ ┌─────────────────┐ │ AI Worker │ 多线程 / 多进程消费 └─────────────────┘ │ │ 调用模型 / 工具 ▼ ┌─────────────────┐ │ AI 服务 │ LLM API / 本地模型 / 脚本 └─────────────────┘ │ └── 处理完成 → 结果回写 Inbox → 客户端查询API 层负责接收用户请求把任务写入 Inbox。Inbox 存储任务记录和状态Worker 定期从中领取任务。领取后Worker 调用 AI 处理服务最后把结果回写到对应任务记录上。线程模型上我们的示例采用「API 线程 后台 Worker 线程」的方式。FastAPI 启动时拉起一个后台 Worker 线程Worker 轮询 Inbox每次领取一个任务。生产环境可以把这个模型替换成多进程 Worker 或独立部署的任务消费者。2.2 任务生命周期任务在收件箱中会经历多个状态。这里把状态定义清楚是整个系统设计的核心。PENDING任务已进入收件箱等待 Worker 领取。PROCESSING任务被某个 Worker 领取正在调用 AI 处理。DONE处理成功结果字段已回填。FAILED处理多次重试仍然失败错误信息已记录。状态流转可以用下面一段伪代码表示提交任务 -- PENDING Worker 领取 -- PROCESSING 处理成功 -- DONE 处理失败且还有重试次数 -- PENDING 处理失败且达到最大次数 -- FAILED这里把“失败后重试”和“失败最终态”区分开非常关键。AI 模型接口经常因为网络抖动、限流、内容审核等原因失败如果一律进入 FAILED会让很多本来可以成功的任务白白失败如果无限重试又会造成成本浪费和队列堆积。常见方案是设置最大尝试次数比如 3 次前 2 次失败回到 PENDING第 3 次失败进入 FAILED。3. 环境准备与项目结构3.1 运行环境与依赖本文示例代码使用 Python 3.10主要依赖 FastAPI 和 Uvicorn。数据库方面先用内存存储演示后续可以在最佳实践章节替换为 Redis 或 SQLite。你需要准备的环境如下Python 3.10 或更高版本。一个虚拟环境venv 或 conda 均可。pip 安装 fastapi、uvicorn、pydantic。版本不需要刻意固定本文示例以常见环境为准重点演示设计思路。下面代码基于 pydantic v2 编写如果你使用的是 pydantic v1需要把model_dump(modejson)改回dict()或者直接在模型中写自定义序列化方法。3.2 项目目录结构为了方便阅读我们将代码拆成几个模块结构如下ai-inbox-assistant/ ├── requirements.txt ├── app/ │ ├── __init__.py │ ├── models.py │ ├── inbox.py │ ├── worker.py │ └── main.py └── README.md各个文件职责如下文件职责requirements.txt项目依赖app/models.py任务数据模型、状态枚举、优先级枚举app/inbox.pyInbox 存储 任务状态管理 领取策略app/worker.py后台消费线程领取任务并调用 AI 服务app/main.pyFastAPI 应用注册路由和生命周期4. 核心模块设计详解4.1 消息模型定义先定义收件箱中的任务模型。它比普通队列消息多了一些业务字段发送方、主题、优先级、内容、状态、尝试次数、结果和错误信息。字段设计说明id任务唯一标识用 UUID 生成。sender消息来源比如email、user、cron。topic任务主题或分类方便后续筛选。content要交给 AI 处理的原始内容。priority优先级影响 Worker 取数的先后顺序。status任务当前状态。attempts当前已尝试执行次数用于失败重试。result处理成功后的结果。error最近一次失败的错误信息。优先级建议使用枚举这样在 API 层校验参数时更安全。我们定义TaskPriority枚举包含HIGH、NORMAL、LOW三档。任务状态用TaskStatus枚举包含PENDING、PROCESSING、DONE、FAILED四种。4.2 Inbox 存储与取数策略这里我们用 Python 内置字典作为任务存储用RLock保证线程安全。收件箱需要提供以下能力add添加任务返回新任务对象。claim_next领取下一个待处理任务。complete完成任务回填结果。fail_or_retry处理失败判断是否重试或进入最终失败态。get按 ID 查询任务。list列出任务支持按状态筛选。claim_next是核心方法。Worker 调用它时必须保证“找出任务”和“修改状态为 PROCESSING”是原子的否则多个 Worker 同时消费时会拿到同一个任务造成重复处理。这里我们在锁内完成查找和状态更新保证了单进程内多个线程不会重复领取。取数策略上我们支持按优先级排序。同一优先级的任务按创建时间先后处理这样既满足业务紧急度要求又不会让低优先级任务无限积压。4.3 Worker 消费端Worker 是一个后台线程循环执行以下步骤调用claim_next()领取任务。如果当前没有任务休眠 1 秒再继续。调用 AI 处理逻辑。成功则调用complete()回填结果。失败则调用fail_or_retry()记录错误并决定是否重试。Worker 使用daemon线程的原因是不阻塞主进程退出。在真实生产环境中建议用进程管理工具或容器编排来管理多个 Worker而不是单线程。4.4 AI 处理服务抽象为了演示我们把“AI 处理”抽象成handle()方法。真实项目中这个方法内部可以调用 OpenAI 等大模型 API也可以调用本地部署的模型服务还可以执行一段工具脚本。这里有一个设计要点AI 调用一定要设置超时。模型接口的响应时间往往不稳定如果 Worker 因为没有超时而卡在一个任务上后续所有任务都会被阻塞。我们可以在handle()中显式设置 HTTP 客户端超时或者用Signal强制中断同步调用。5. 完整实现FastAPI 线程 Worker5.1 创建项目与安装依赖首先创建项目目录和虚拟环境。mkdir ai-inbox-assistant cd ai-inbox-assistant python -m venv .venv source .venv/bin/activate # Windows 使用 .venv\Scripts\activate创建requirements.txt并写入以下依赖fastapi0.110 uvicorn[standard]0.29 pydantic2.0安装依赖pip install -r requirements.txt5.2 定义数据模型文件路径app/models.pyfrom datetime import datetime, timezone from enum import Enum from typing import Optional from pydantic import BaseModel, Field def utc_now() - datetime: 统一获取当前 UTC 时间避免重复实现。 return datetime.now(timezone.utc) class TaskStatus(str, Enum): PENDING pending PROCESSING processing DONE done FAILED failed class TaskPriority(str, Enum): HIGH high NORMAL normal LOW low class InboxTask(BaseModel): id: str Field(default_factorylambda: __import__(uuid).uuid4().hex) sender: str unknown topic: str default content: str priority: TaskPriority TaskPriority.NORMAL status: TaskStatus TaskStatus.PENDING created_at: datetime Field(default_factoryutc_now) updated_at: datetime Field(default_factoryutc_now) attempts: int 0 result: Optional[str] None error: Optional[str] None def to_dict(self) - dict: 转为可直接 JSON 序列化的字典兼容 pydantic v1/v2。 return { id: self.id, sender: self.sender, topic: self.topic, content: self.content, priority: self.priority.value, status: self.status.value, created_at: self.created_at.isoformat(), updated_at: self.updated_at.isoformat(), attempts: self.attempts, result: self.result, error: self.error, }这里重点解释几个设计决策。id使用uuid4().hex生成 32 位十六进制字符串足以避免并发提交时的 ID 冲突。也可以直接用str(uuid.uuid4())区别只是是否带横线。to_dict()方法统一负责序列化把枚举值、时间对象转换为普通字符串这样接口层在返回响应时不需要关心底层 pydantic 版本差异。attempts字段默认 0表示任务还未被消费。5.3 实现 Inbox 核心逻辑文件路径app/inbox.pyimport threading from typing import Dict, List, Optional from .models import InboxTask, TaskPriority, TaskStatus class Inbox: 线程安全的内存收件箱。 def __init__(self) - None: self._tasks: Dict[str, InboxTask] {} self._lock threading.RLock() staticmethod def _sort_key(task: InboxTask): 优先级高的任务排在前面相同优先级按创建时间判断。 priority_order { TaskPriority.HIGH: 0, TaskPriority.NORMAL: 1, TaskPriority.LOW: 2, } return (priority_order.get(task.priority, 1), task.created_at) def add(self, content: str, sender: str unknown, topic: str default, priority: TaskPriority TaskPriority.NORMAL) - InboxTask: 向收件箱添加一个任务。 with self._lock: task InboxTask( sendersender, topictopic, contentcontent, prioritypriority, ) self._tasks[task.id] task return task def claim_next(self) - Optional[InboxTask]: 领取下一个待处理任务并将状态改为 PROCESSING。 with self._lock: candidates [ task for task in self._tasks.values() if task.status TaskStatus.PENDING ] if not candidates: return None candidates.sort(keyself._sort_key) task candidates[0] task.status TaskStatus.PROCESSING task.attempts 1 task.updated_at __import__(app.models, fromlist[utc_now]).utc_now() return task def complete(self, task_id: str, result: str) - None: 处理成功后回填结果。 with self._lock: task self._tasks.get(task_id) if task is None: raise KeyError(ftask {task_id} not found) task.status TaskStatus.DONE task.result result task.error None task.updated_at __import__(app.models, fromlist[utc_now]).utc_now() def fail_or_retry(self, task_id: str, error: str, max_attempts: int 3) - None: 错误处理如果未超过最大执行次数则回到 PENDING否则标记 FAILED。 with self._lock: task self._tasks.get(task_id) if task is None: raise KeyError(ftask {task_id} not found) task.error error task.updated_at __import__(app.models, fromlist[utc_now]).utc_now() if task.attempts max_attempts: task.status TaskStatus.PENDING else: task.status TaskStatus.FAILED def get(self, task_id: str) - Optional[InboxTask]: with self._lock: return self._tasks.get(task_id) def list(self, status: Optional[TaskStatus] None) - List[InboxTask]: with self._lock: tasks list(self._tasks.values()) if status is not None: tasks [t for t in tasks if t.status status] tasks.sort(keylambda t: t.created_at) return tasks这里有一个实现细节需要注意claim_next返回的是任务对象本身而不是副本。这意味着 Worker 在拿到任务对象后即使 Inbox 锁已经释放其他线程读取这个任务时也能看到PROCESSING状态。这正是我们希望的效果因为它反映了真实的执行状态。不过要注意由于我们直接修改任务对象的属性如果 Worker 在执行任务时不小心修改了content等业务字段会产生脏数据。所以在设计约定上Worker 只允许通过complete和fail_or_retry修改任务状态不要直接操作字段。5.4 实现 Worker 消费端文件路径app/worker.pyimport threading import time from typing import Optional from .inbox import Inbox from .models import InboxTask class AIWorker(threading.Thread): 后台消费线程从 Inbox 领取任务、调用模型、回填结果。 def __init__(self, inbox: Inbox, name: str ai-worker, poll_interval: float 1.0): super().__init__(namename, daemonTrue) self.inbox inbox self.poll_interval poll_interval self._stop_event threading.Event() def stop(self) - None: self._stop_event.set() def run(self) - None: while not self._stop_event.is_set(): task: Optional[InboxTask] self.inbox.claim_next() if task is None: self._stop_event.wait(self.poll_interval) continue try: result self.handle(task) self.inbox.complete(task.id, result) except Exception as exc: self.inbox.fail_or_retry(task.id, str(exc)) def handle(self, task: InboxTask) - str: 核心 AI 处理函数可替换为真实模型 API 调用。 # 模拟耗时操作生产环境替换成 LLM API / 本地模型推理 time.sleep(0.5) return f[{task.topic}] {task.content[:20]} 的 AI 摘要已生成Worker 中最容易踩坑的是异常处理边界。handle()中任何异常都会触发fail_or_retry()但这个逻辑需要与重试策略配合。如果任务是“永久性错误”例如内容包含非法字符导致模型拒绝处理重试多少次都不成功反而会浪费资源。所以在生产系统中handle()内部应该区分临时错误和永久错误永久错误直接抛出特定异常由调用方判断是一次性失败还是继续重试。这个示例中的poll_interval是 1 秒在演示环境可以接受。生产环境通常用消息队列的阻塞读取或者长轮询避免无意义的轮询开销。5.5 编写 FastAPI 接口文件路径app/main.pyfrom contextlib import asynccontextmanager from typing import Optional from fastapi import FastAPI, HTTPException, Query from .inbox import Inbox from .models import InboxTask, TaskStatus from .worker import AIWorker inbox Inbox() worker: Optional[AIWorker] None asynccontextmanager async def lifespan(app: FastAPI): global worker worker AIWorker(inbox, nameai-worker) worker.start() yield if worker is not None: worker.stop() app FastAPI( titleAI Assistant with Inbox, description一个自带收件箱的 AI 助手服务, version0.1.0, lifespanlifespan, ) class TaskCreateRequest: def __init__(self, content: str, sender: str unknown, topic: str default, priority: str normal): self.content content self.sender sender self.topic topic self.priority priority from pydantic import BaseModel class TaskCreateBody(BaseModel): content: str sender: str unknown topic: str default priority: str normal class TaskListResponse(BaseModel): items: list[dict] app.post(/inbox/tasks, status_code201) def create_task(body: TaskCreateBody) - dict: 提交一个新任务到收件箱。 from .models import TaskPriority try: priority TaskPriority(body.priority) except ValueError: raise HTTPException(status_code422, detailf无效优先级: {body.priority}) task inbox.add( contentbody.content, senderbody.sender, topicbody.topic, prioritypriority, ) return {task_id: task.id, status: task.status.value} app.get(/inbox/tasks) def list_tasks( status: Optional[TaskStatus] Query(defaultNone), sender: Optional[str] Query(defaultNone), ) - TaskListResponse: 列出收件箱任务支持按状态和发送方筛选。 tasks inbox.list(statusstatus) if sender: tasks [t for t in tasks if t.sender sender] return TaskListResponse(items[t.to_dict() for t in tasks]) app.get(/inbox/tasks/{task_id}) def get_task(task_id: str) - dict: 查询单个任务状态和结果。 task: Optional[InboxTask] inbox.get(task_id) if task is None: raise HTTPException(status_code404, detail任务不存在) return task.to_dict()代码里保留了TaskCreateRequest这个旧类其实是不需要的可以去掉。我在这里故意保留是因为实际开发中经常会有“写多了再清理”的情况。正式代码建议直接删掉只保留 Pydantic 模型。接口设计有三个核心点。第一POST /inbox/tasks返回task_id而不是完整处理结果。客户端拿到任务 ID 后可以通过GET /inbox/tasks/{task_id}轮询结果。这是异步任务接口的标准做法。第二查询接口支持按status和sender过滤方便业务侧按状态或来源查看收件箱内容。第三状态枚举通过 Query 参数接收时FastAPI 会自动做参数校验。如果传入非法状态返回 422不需要我们手写校验逻辑。5.6 启动服务并验证现在启动服务。uvicorn app.main:app --reload --port 8000看到如下输出说明启动成功INFO: Uvicorn running on http://127.0.0.1:8000 INFO: Application startup complete.FastAPI 会自动生成交互式文档访问http://127.0.0.1:8000/docs可以查看所有接口。6. 运行演示与结果说明6.1 提交任务打开另一个终端使用 curl 提交两个测试任务一个高优先级一个普通优先级。curl -X POST http://127.0.0.1:8000/inbox/tasks \ -H Content-Type: application/json \ -d {content: 请总结本周运营数据, sender: cron, topic: report, priority: high}预期输出{task_id:9f7b2f6d0c9a4e6f9c48e0c9ae62da21,status:pending}再提交一个普通任务curl -X POST http://127.0.0.1:8000/inbox/tasks \ -H Content-Type: application/json \ -d {content: 生成一封客户回复邮件, sender: user, topic: email, priority: normal}6.2 查询任务列表任务提交后立即查询列表可能看到部分任务处于pending部分处于processing取决于 Worker 的处理速度。curl http://127.0.0.1:8000/inbox/tasks输出示例{ items: [ { id: 9f7b2f6d0c9a4e6f9c48e0c9ae62da21, sender: cron, topic: report, content: 请总结本周运营数据, priority: high, status: done, created_at: 2025-01-01T10:00:0000:00, updated_at: 2025-01-01T10:00:0100:00, attempts: 1, result: [report] 请总结本周运营数据 的 AI 摘要已生成, error: null } ] }注意attempts字段已经变成 1说明 Worker 领取并处理过一次。result字段已经回填生成结果。6.3 查询单个任务结果根据之前拿到的task_id查询单个任务curl http://127.0.0.1:8000/inbox/tasks/9f7b2f6d0c9a4e6f9c48e0c9ae62da21输出与列表中的单个项目一致。到这里一个最小的“带收件箱的 AI 助手”已经可以跑通了。7. 常见问题与排查思路实际开发中你会遇到各种预期外的情况。下面整理了一些高频问题。问题现象常见原因解决思路任务一直 pending状态不变Worker 线程没有启动或提前退出检查 lifespan 是否生效打印 Worker 启动日志多个 Worker 重复处理同一任务领取任务和修改状态不是原子操作在锁/事务中完成状态更新使用分布式锁任务失败后频繁重试没有区分临时错误和永久错误定义可重试异常永久错误直接标记 FAILED服务重启后任务丢失任务存储在内存中引入 Redis Streams、SQLite、PostgreSQL 持久化API 返回 422状态或优先级参数传错核对枚举值大小写参考 /docs 接口文档模型调用超时导致 Worker 卡死外部接口没有设置超时为 AI 调用设置超时时间并配合重试策略Uvicorn 启动报 lifespan 错误代码缩进或局部变量问题检查 lifespan 上下文管理器结构启动日志会显示堆栈7.1 任务一直处于 pending 状态出现这个现象首先检查 Worker 是否在运行。在启动日志中看不到 Worker 相关信息时往往是 lifespan 生命周期没有挂载正确。FastAPI 旧版本常见做法是app.on_event(startup)新版开始推荐lifespan上下文管理器。如果你使用的是较老版本 FastAPI可以改回 startup 事件写法但要注意不同版本的兼容性。还可以在 Worker 的run()方法最开始加一行打印日志比如print([worker] started)这样能很快确认线程是否启动。7.2 任务重复消费在单进程多线程模型中claim_next因为有RLock保护不会出现重复领取。但在多进程部署时每个进程都有自己的 Inbox 实例任务存储不共享这时候问题会变成“各进程各处理各的”而不是重复消费同一个任务。真正的重复消费风险发生在任务存储是共享的比如 Redis但领取时没有用原子操作。解决方法有两种在 Inbox 存储层使用带条件的原子更新例如 Redis Lua 脚本或 SQLUPDATE ... WHERE statuspending。在 Worker 处理结果回写时使用幂等 ID 校验防止重复写入结果。对于 AI 任务重复消费不只是资源浪费还可能导致重复扣费和重复生成内容所以幂等设计要提前做。7.3 模型调用超时模型 API 是外部依赖它的延迟不可控。如果不设置超时一个慢请求可能让 Worker 长期阻塞。常见做法有在网络请求库层面设置timeout比如requests.post(url, timeout(3, 30))。在多线程 Worker 中用Future.get(timeout...)控制单个任务执行时长。为任务设置最大执行时间超过阈值的任务重新进入队列或直接标记失败。8. 最佳实践与工程建议演示代码跑通后如果要在生产环境落地下面这些点非常关键。8.1 存储层选型内存字典最明显的缺点是重启丢数据。生产环境推荐替换为以下方案之一。存储方案适合场景优点注意点Redis Streams中高吞吐任务队列天然支持消息持久化、消费者组需要处理 Stream 的消息过期和积压Redis List BRPOP简单任务队列实现简单阻塞读取缺少消费者 ACK需要额外设计SQLite 状态列低并发单机任务零额外依赖方便审计写并发有限需要适当加锁PostgreSQL SKIP LOCKED中大型系统支持事务和 SKIP LOCKED 避免重复消费需要数据库连接池如果你已经有 RabbitMQ 或 Kafka 基础设施也可以直接把它们作为任务队列但要在消息体里保留task_id和完整错误信息。8.2 幂等与重试策略AI 调用通常涉及成本重试策略必须谨慎。建议按以下原则设计为每个任务生成全局唯一request_id发往模型服务时携带该 ID。网络超时、限流、5xx 等临时错误允许重试。内容不合法、参数错误等永久错误不要重试。设置最大尝试次数默认为 3避免无限重试。使用指数退避策略比如第 1 次等 2 秒第 2 次等 4 秒第 3 次等 8 秒。在当前的fail_or_retry方法中最简单的指数退避可以放在 Worker 内部实现重试前time.sleep(backoff)。8.3 超时与死信任务长时间处于PROCESSING状态可能是 Worker 崩溃导致的任务“死亡”。生产环境需要引入“死信”机制。可以每隔一段时间扫描状态为PROCESSING但updated_at超过 10 分钟的任务将它们重新置为PENDING或标记为FAILED并记录告警。这个扫描任务通常由定时调度器执行。8.4 安全与鉴权收件箱中可能包含敏感数据比如客户邮件、业务报告文本。接口不能裸奔在公网上。建议在 FastAPI 中配置 API Key 或 OAuth2 鉴权。对任务内容加密存储。查询接口做权限校验普通用户只能查询自己提交的任务不能查看他人的任务内容。记录每个请求的操作人、时间和任务 ID以便审计。8.5 AI 调用成本控制当收件箱堆积大量任务时如果不做控制模型 API 账单会很快飙升。控制成本可以从几个方向入手任务入库前进行内容长度限制和去重。对相同或近似内容做缓存命中后直接返回历史结果。给 Worker 加速率限制防止瞬间请求过多导致模型 API 限流。流式读取大文本时先做预处理
返回列表
PREV
查看更多资讯
NEXT
返回资讯列表