Skip to content

第 3 章 · 分布式协调与容器化部署进阶

本章目标:

  • 掌握 Redis 分布式锁的正确实现方式,避免重复消费
  • 理解幂等性设计在 Agent 任务中的重要性
  • 学会使用 Docker 多阶段构建优化 Python 镜像大小
  • 能用 docker-compose 编排多容器 Agent 服务
  • 为生产环境部署打下容器化基础

在单机上跑的 Agent 是玩具,能在多节点上可靠协作的 Agent 才是生产系统。本章将学习如何在分布式环境中协调多个 Worker,并用 Docker 把这些组件打包成可移植的服务。

3.1 分布式锁:防止重复消费

当多个 Celery Worker 同时监听同一个队列时,可能出现「一个任务被多个 Worker 抢走」的情况。用 Redis 的 SET NX EX 实现原子性互斥锁:

python
# services/agent/redis_lock.py
import redis
import uuid
from contextlib import contextmanager

r = redis.from_url("redis://redis:6379/0")


@contextmanager
def distributed_lock(lock_name: str, ttl: int = 30):
    """
    用 Redis SET NX EX 实现可重入分布式锁。

    - NX:键不存在才设置(互斥)
    - EX:过期时间(防止 Worker 崩溃后锁永不过期)
    - 唯一 token:保证只有持锁者能释放锁
    """
    token = str(uuid.uuid4())
    acquired = r.set(lock_name, token, nx=True, ex=ttl)
    if not acquired:
        raise RuntimeError(f"无法获取锁: {lock_name}")
    try:
        yield
    finally:
        # Lua 脚本保证「检查 token + 删除」的原子性
        lua = """
        if redis.call("get", KEYS[1]) == ARGV[1] then
            return redis.call("del", KEYS[1])
        else
            return 0
        end
        """
        r.eval(lua, 1, lock_name, token)

在 Celery 任务中使用

python
# services/agent/tasks.py
from celery import shared_task
from .redis_lock import distributed_lock


@shared_task(bind=True, max_retries=3)
def process_agent_task(self, task_id: str, user_input: str):
    """
    用分布式锁确保同一 task_id 在同一时刻只被一个 Worker 处理。
    """
    lock_name = f"agent:task:{task_id}"
    with distributed_lock(lock_name, ttl=60):
        # 业务逻辑:调用 LangGraph 图
        result = run_agent_graph(user_input)
        return {"task_id": task_id, "result": result}

为什么用 Lua 脚本释放锁?

如果只用 GET + DEL 两步,存在竞态条件:Worker A 拿到锁后崩溃,Worker B 在过期前检查到锁存在就会跳过,直到 TTL 到期。而「检查 token + 删除」用 Lua 脚本保证原子性,避免误删他人持有的锁。

3.2 幂等性设计:唯一任务 ID 与去重

分布式环境下任务可能重复投递(网络抖动、Worker 重试)。解决方案:

  1. 客户端生成唯一 task_id(UUID v7,含时间戳可排序)
  2. 任务开始前查去重表,已处理则直接返回缓存结果
  3. 任务完成后写入结果表
python
# services/agent/idempotency.py
import uuid
from datetime import datetime
import redis

r = redis.from_url("redis://redis:6379/0")


def generate_task_id() -> str:
    """生成 UUID v7,便于时间排序和索引。"""
    return str(uuid.uuid7())


def check_duplicate(task_id: str) -> bool:
    """
    检查 task_id 是否已被处理过。
    Redis SET NX 原子操作:首次返回 True(未处理),重复返回 False。
    """
    key = f"idempotent:{task_id}"
    return bool(r.set(key, "1", nx=True, ex=3600))  # 1 小时过期


def cache_result(task_id: str, result: dict) -> None:
    """缓存任务结果,供重复请求直接返回。"""
    r.setex(f"result:{task_id}", 300, str(result))  # 5 分钟 TTL
python
# services/agent/tasks.py
from .idempotency import generate_task_id, check_duplicate, cache_result


@shared_task
def handle_agent_request(user_input: str):
    task_id = generate_task_id()

    if not check_duplicate(task_id):
        return {"task_id": task_id, "status": "duplicate", "note": "任务已处理"}

    # 执行 Agent 逻辑...
    result = {"response": "Agent 输出...", "task_id": task_id}
    cache_result(task_id, result)
    return result

💡 生产建议:高频场景下用 Postgres 去重表替代 Redis,利用 INSERT ... ON CONFLICT DO NOTHING 实现强一致去重。详见 FastAPI 课程第 13 章

3.3 Docker 多阶段构建:优化 Python 镜像

Python 镜像默认带完整 GCC 和开发头文件,体积可达 1GB。多阶段构建可以大幅缩小:

dockerfile
# services/agent/Dockerfile
FROM python:3.12-slim AS base

# 阶段 1:安装系统依赖
FROM base AS builder
RUN apt-get update && apt-get install -y --no-install-recommends \
    gcc libpq-dev && \
    rm -rf /var/lib/apt/lists/*

# 复制依赖文件并安装 Python 包
COPY requirements.txt .
RUN pip install --no-cache-dir --prefix=/install -r requirements.txt

# 阶段 2:精简最终镜像
FROM base
WORKDIR /app
COPY --from=builder /install /usr/local
COPY . .

# 非 root 用户运行
RUN useradd -m agentuser && chown -R agentuser /app
USER agentuser

EXPOSE 8000
CMD ["uvicorn", "services.agent.main:app", "--host", "0.0.0.0", "--port", "8000"]

关键优化点

优化项效果
python:3.12-slim基础镜像从 1GB 降到 ~150MB
--no-install-recommends避免安装不必要的推荐包
--no-cache-dirpip 不缓存wheel,减少镜像体积
非 root 用户安全最佳实践,防止容器逃逸

3.4 docker-compose.yml:编排多容器服务

一个生产 Agent 服务通常包含:FastAPI 网关、Celery Worker、Redis、Postgres。用 docker-compose.yml 一键启动:

yaml
# docker-compose.yml
services:
  api:
    build:
      context: .
      dockerfile: services/agent/Dockerfile
    ports:
      - "8000:8000"
    environment:
      - REDIS_URL=redis://redis:6379/0
      - DATABASE_URL=postgresql://agent:agent@postgres:5432/agentdb
      - CELERY_BROKER_URL=redis://redis:6379/1
      - CELERY_RESULT_BACKEND=redis://redis:6379/2
    depends_on:
      redis: { condition: service_healthy }
      postgres: { condition: service_healthy }
    healthcheck:
      test: ["CMD", "curl", "-f", "http://localhost:8000/health"]
      interval: 30s
      timeout: 10s
      retries: 3

  worker:
    build:
      context: .
      dockerfile: services/agent/Dockerfile
    command: celery -A services.agent.celery_app worker --loglevel=info
    environment:
      - REDIS_URL=redis://redis:6379/0
      - DATABASE_URL=postgresql://agent:agent@postgres:5432/agentdb
      - CELERY_BROKER_URL=redis://redis:6379/1
    depends_on:
      redis: { condition: service_healthy }
      postgres: { condition: service_healthy }

  redis:
    image: redis:7-alpine
    ports:
      - "6379:6379"
    volumes:
      - redis_data:/data
    healthcheck:
      test: ["CMD", "redis-cli", "ping"]
      interval: 10s
      timeout: 5s
      retries: 5

  postgres:
    image: postgres:16-alpine
    environment:
      POSTGRES_USER: agent
      POSTGRES_PASSWORD: agent
      POSTGRES_DB: agentdb
    ports:
      - "5432:5432"
    volumes:
      - postgres_data:/var/lib/postgresql/data
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U agent"]
      interval: 10s
      timeout: 5s
      retries: 5

volumes:
  redis_data:
  postgres_data:

📌 健康检查depends_on.condition: service_healthy 确保依赖服务真正可用后才启动,避免竞态。

3.5 健康检查与优雅关闭

FastAPI 健康端点

python
# services/agent/main.py
from fastapi import FastAPI
from redis import Redis
import psycopg2

app = FastAPI()


@app.get("/health")
async def health_check():
    checks = {}

    # Redis 连通性
    try:
        r = Redis.from_url("redis://redis:6379/0")
        checks["redis"] = r.ping()
    except Exception as e:
        checks["redis"] = False

    # Postgres 连通性
    try:
        conn = psycopg2.connect("postgresql://agent:agent@postgres:5432/agentdb")
        conn.close()
        checks["postgres"] = True
    except Exception:
        checks["postgres"] = False

    healthy = all(checks.values())
    status_code = 200 if healthy else 503
    return {"status": "healthy" if healthy else "unhealthy", "checks": checks}, status_code

Celery Worker 优雅关闭

python
# services/agent/celery_app.py
from celery import Celery
import signal
import sys

app = Celery("agent", broker="redis://redis:6379/1")
app.conf.update(
    task_serializer="json",
    result_serializer="json",
    accept_content=["json"],
    timezone="UTC",
    enable_utc=True,
)


def graceful_shutdown(signum, frame):
    print("收到退出信号,完成当前任务后停止...")
    sys.exit(0)


signal.signal(signal.SIGTERM, graceful_shutdown)
signal.signal(signal.SIGINT, graceful_shutdown)

💡 Docker 信号传递:在 Dockerfile 中添加 ENTRYPOINT ["tini", "--"] 确保进程能正确接收 SIGTERM 信号。

3.6 完整示例:多容器 Agent 服务

项目结构

agentic-tutorial/
├── services/
│   └── agent/
│       ├── Dockerfile
│       ├── main.py
│       ├── celery_app.py
│       ├── tasks.py
│       ├── redis_lock.py
│       ├── idempotency.py
│       └── requirements.txt
├── docker-compose.yml
└── README.md

requirements.txt

txt
fastapi==0.141.0
uvicorn[standard]==0.34.0
celery[redis]==5.4.0
redis==5.2.1
psycopg2-binary==2.9.10
langgraph==1.2.11
anthropic==0.49.0
pydantic-settings==2.7.1

启动命令

bash
# 后台启动所有服务
docker-compose up -d

# 查看日志
docker-compose logs -f api worker

# 停止服务
docker-compose down

本章小结

  • 分布式锁:用 SET NX EX + Lua 脚本保证原子性,防止重复消费
  • 幂等性:UUID v7 + Redis 去重表,应对网络重试和 Worker 重启
  • 多阶段构建:slim 基础镜像 + 分离 builder 阶段,镜像体积缩小 80%
  • docker-composecondition: service_healthy 确保依赖顺序
  • 健康检查/health 端点 + HEALTHCHECK 指令,支持 Kubernetes 探针

⏭️ 下一章将进入 LangGraph 核心:StateGraph 与工具调用,这是构建生产 Agent 的基石。

🛠️ 动手实践

  1. 实现 Redis 分布式锁:在 services/agent/redis_lock.py 基础上,增加可重入支持——同一 Worker 多次获取同一锁不阻塞,释放次数与获取次数匹配。
  2. 优化 Dockerfile:将当前 1GB 左右的 Python 镜像压缩到 200MB 以内,使用 python:3.12-slim + 多阶段构建,并添加 .dockerignore 排除 .gitnode_modules
  3. 编写 healthcheck:为 docker-compose.yml 中的 api 服务添加 /health 端点,要求 Redis 和 Postgres 都健康才返回 200,任一失败返回 503。

🧪 随堂测验

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

1. Redis SET NX EX 命令中,NX 参数的作用是什么?

2. 以下哪种方式可以保证锁释放的原子性?

3. Docker 多阶段构建中,builder 阶段的主要作用是什么?

4. docker-compose.yml 中 depends_on 的 condition: service_healthy 作用是什么?