Skip to content

第 22 章 · 部署与服务化

本章目标:对比 CrewAI Enterprise 托管与自托管两条路线,用 FastAPI 把 crew 包装成异步任务 API,掌握环境变量注入、长任务的 checkpoint 恢复与无状态水平扩展。

22.1 两条部署路线怎么选

维度CrewAI Enterprise(托管)自托管(本节重点)
上手成本crewai deploy 推送即上线自己写服务层、管进程
运行环境平台托管、自动伸缩Docker/K8s 自管
数据边界运行在平台侧完全自有
灵活性受平台约束任意改造

学习路径建议

先自托管跑通"API 化 → 超时恢复 → 扩展",理解全部运维关注点后,再评估是否迁移到 Enterprise 省运维。两者概念一一对应:部署 = 服务化,触发 = kickoff,日志 = 第 21 章的观测体系。

自托管的最小形态就是一个 Web 服务接收输入、执行 crew.kickoff()、返回结果。

22.2 用 FastAPI 包装 kickoff:异步任务 + 轮询

生产中一次 kickoff 可能跑几分钟,HTTP 同步等待会超时。标准模式是 POST 提交任务返回 job_id → 后台执行 → GET 轮询状态

python
# main.py
import os
import uuid
import asyncio
from fastapi import FastAPI, HTTPException, BackgroundTasks

from marketing_crew import MarketingCrew  # 你的 crew 定义(第 23 章)

app = FastAPI(title="Marketing Crew API")

# 进程内任务表;多副本部署时换成 Redis/数据库(见 22.4)
JOBS: dict[str, dict] = {}


def run_crew_job(job_id: str, topic: str) -> None:
    """在线程池中同步执行 crew,并把结果写回任务表"""
    try:
        result = MarketingCrew().crew().kickoff(inputs={"topic": topic})
        JOBS[job_id].update(
            status="completed",
            result=result.raw,
            tokens=result.token_usage.total_tokens,
        )
    except Exception as e:
        JOBS[job_id].update(status="failed", error=str(e))


@app.post("/kickoff")
def kickoff(topic: str, background_tasks: BackgroundTasks):
    """提交一次 crew 运行,立即返回 job_id"""
    if not os.getenv("DEEPSEEK_API_KEY"):
        raise HTTPException(500, "缺少 DEEPSEEK_API_KEY 环境变量")
    job_id = uuid.uuid4().hex
    JOBS[job_id] = {"status": "running", "result": None}
    # FastAPI BackgroundTasks 在响应返回后执行函数
    background_tasks.add_task(run_crew_job, job_id, topic)
    return {"job_id": job_id, "status": "running"}


@app.get("/status/{job_id}")
def status(job_id: str):
    """客户端轮询任务状态"""
    job = JOBS.get(job_id)
    if not job:
        raise HTTPException(404, "任务不存在")
    return job


# 启动: uvicorn main:app --host 0.0.0.0 --port 8000
# curl -X POST "http://localhost:8000/kickoff?topic=AI+Agent"
# curl http://localhost:8000/status/<job_id>

要点:BackgroundTasks 让 HTTP 立即返回;crew 本身是同步阻塞的,放进后台线程而不是 async def 里直接调(会卡死事件循环)。若要更强的并发控制,改用 asyncio.to_thread 或独立任务队列(Celery/RQ/Arq)。

22.3 密钥与环境变量注入

任何情况下都不要把 API key 写进代码或镜像:

dockerfile
# Dockerfile —— 多阶段构建,密钥只通过运行时注入
FROM python:3.12-slim AS base
WORKDIR /app

COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

COPY . .

# 不在镜像里放任何密钥!
ENV PYTHONUNBUFFERED=1
EXPOSE 8000
CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8000"]
bash
# 运行时通过 env 文件注入(.env 加入 .gitignore)
docker run --rm -p 8000:8000 \
  --env-file .env \
  my-crew-api:latest

# .env 内容示例:
# DEEPSEEK_API_KEY=sk-xxx
# SERPER_API_KEY=xxx

CrewAI 会从进程环境读取 *_API_KEY 等变量;Kubernetes 场景用 Secret 挂载为环境变量,效果相同。

22.4 长任务的超时与 checkpoint 恢复

跑到一半被杀(OOM、发布重启、上游 LLM 抖动)怎么办?CrewAI 内置检查点机制:按事件保存可恢复的完整状态快照,恢复时自动跳过已完成的任务:

python
from crewai import Agent, Task, Crew

researcher = Agent(role="调研员", goal="调研行业趋势", backstory="资深分析师")
writer = Agent(role="撰稿人", goal="撰写报告", backstory="技术作家")

crew = Crew(
    agents=[researcher, writer],
    tasks=[
        Task(description="调研 AI Agent 市场", agent=researcher,
             expected_output="要点列表"),
        Task(description="撰写总结报告", agent=writer,
             expected_output="500 字报告"),
    ],
    checkpoint=True,  # 开启事件驱动检查点,默认每个 task_completed 存一份到 ./.checkpoints/
)

result = crew.kickoff()  # 中途 Ctrl+C 或进程被杀后:
# ./.checkpoints/<timestamp>_<uuid>.json 就是快照
python
from pathlib import Path

# 从最新快照恢复:已完成的任务会被跳过,直接续跑剩余部分
checkpoint_file = sorted(Path("./.checkpoints").glob("*.json"))[-1]
resumed = Crew(
    agents=[researcher, writer],
    tasks=crew.tasks,
    checkpoint=checkpoint_file,   # 传入快照文件即进入恢复模式
)
print(resumed.kickoff().raw)

配套参数:存储提供 JsonProvider(一文件一快照,易检视)与 SqliteProvider(单库高频写入);max_checkpoints 控制保留数量、自动淘汰最旧快照。注意自动检查点是尽力而为——写入失败仅记日志不中断运行;手动 state.checkpoint() 则失败即抛错。把 .checkpoints/ 目录挂成持久卷,配合 22.2 的 job 表就能实现"服务重启后继续跑未完成任务"。

22.5 水平扩展:让服务无状态

多副本部署前先做无状态化审计,四个必查项:

  1. 任务表外置:22.2 的 JOBS 字典是进程内存态——多副本下轮询会打到没有该任务的实例上。换成 Redis(SET job:<id> ... EX 3600)或数据库表;
  2. 文件产物走对象存储output_file 写本地盘的产物,多副本间不可见,改传 S3/OSS;
  3. 记忆/知识库外置:开启 memory/knowledge 时配置共享后端(如外部 embedder 与存储),避免各副本各存一份;
  4. 检查点目录共享.checkpoints/ 挂网络卷或改用 SqliteProvider 放共享库,否则恢复只能落在原副本。

扩容本身很简单——crew 执行是 CPU+IO 型进程,uvicorn --workers N 或 K8s HPA 按 CPU/队列长度伸缩即可;真正的瓶颈通常在上游 LLM 的速率限制(记得用 max_rpm 保护,见第 25 章)。

任务表外置的 Redis 参考实现:

python
# jobs.py —— 用 Redis 替换进程内字典,支持多副本部署
import os
import redis

r = redis.Redis.from_url(os.getenv("REDIS_URL", "redis://localhost:6379"))
TTL = 3600  # 任务记录保留 1 小时


def create_job(job_id: str) -> None:
    r.hset(f"job:{job_id}", mapping={"status": "running", "result": ""})
    r.expire(f"job:{job_id}", TTL)


def finish_job(job_id: str, result: str, tokens: int) -> None:
    r.hset(f"job:{job_id}", mapping={
        "status": "completed", "result": result, "tokens": tokens})


def get_job(job_id: str) -> dict | None:
    data = r.hgetall(f"job:{job_id}")
    if not data:
        return None
    return {k.decode(): v.decode() for k, v in data.items()}

把 22.2 中对 JOBS 字典的读写全部换成这三个函数后,任意副本的 /status 都能查到同一任务——这就是无状态化的核心一步。

本章小结

  • 部署两路线:Enterprise 托管省运维,自托管灵活可控;先自托管吃透运维点;
  • API 化标准姿势:POST /kickoff 返回 job_id + BackgroundTasks 后台执行 + GET /status 轮询;
  • 密钥只在运行时经环境变量注入,镜像构建阶段零密钥;
  • checkpoint=True 事件驱动存快照,恢复时跳过已完成任务;JsonProvider 易检视、SqliteProvider 抗高频;
  • 无状态四查:任务表、文件产物、记忆知识库、检查点目录全部外置后才好水平扩展。

🧪 随堂测验

点击你认为正确的选项。答错时会展示正确答案与原因解析。

1. 为什么 /kickoff 接口要用「后台任务 + job_id 轮询」而不是同步等 kickoff 完成?

2. 关于 CrewAI checkpoint,下列说法错误的是?

3. 把 22.2 的服务扩到 3 个副本后,GET /status 经常 404,最可能的原因是?

4. 关于容器中的密钥管理,正确做法是?

🛠️ 动手实践

  1. 给 22.2 的服务增加 DELETE /jobs/{job_id} 取消接口:把任务标记为 cancelled,并在 crew 执行结束后忽略迟到的结果写入。
  2. 为你的 crew 开启 checkpoint=True,手动在第 2 个任务执行期间 kill -9 进程,然后用最新快照恢复运行,验证第 1 个任务确实被跳过(在其回调里打印日志确认)。
  3. JOBS 字典改造为 Redis 实现(redis-py 的 get/set + TTL),用 docker compose 起 2 个 API 副本 + 1 个 Redis,验证任意副本都能查询到同一 job 的状态。

原理全部就绪,接下来用两个完整项目把它们串起来:第 23 章 · 综合实战一