第 3 章 · 分布式协调与容器化部署进阶
本章目标:
- 掌握 Redis 分布式锁的正确实现方式,避免重复消费
- 理解幂等性设计在 Agent 任务中的重要性
- 学会使用 Docker 多阶段构建优化 Python 镜像大小
- 能用 docker-compose 编排多容器 Agent 服务
- 为生产环境部署打下容器化基础
在单机上跑的 Agent 是玩具,能在多节点上可靠协作的 Agent 才是生产系统。本章将学习如何在分布式环境中协调多个 Worker,并用 Docker 把这些组件打包成可移植的服务。
3.1 分布式锁:防止重复消费
当多个 Celery Worker 同时监听同一个队列时,可能出现「一个任务被多个 Worker 抢走」的情况。用 Redis 的 SET NX EX 实现原子性互斥锁:
# 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 任务中使用
# 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 重试)。解决方案:
- 客户端生成唯一 task_id(UUID v7,含时间戳可排序)
- 任务开始前查去重表,已处理则直接返回缓存结果
- 任务完成后写入结果表
# 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# 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。多阶段构建可以大幅缩小:
# 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-dir | pip 不缓存wheel,减少镜像体积 |
| 非 root 用户 | 安全最佳实践,防止容器逃逸 |
3.4 docker-compose.yml:编排多容器服务
一个生产 Agent 服务通常包含:FastAPI 网关、Celery Worker、Redis、Postgres。用 docker-compose.yml 一键启动:
# 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 健康端点
# 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_codeCelery Worker 优雅关闭
# 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.mdrequirements.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启动命令
# 后台启动所有服务
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-compose:
condition: service_healthy确保依赖顺序 - 健康检查:
/health端点 +HEALTHCHECK指令,支持 Kubernetes 探针
⏭️ 下一章将进入 LangGraph 核心:StateGraph 与工具调用,这是构建生产 Agent 的基石。
🛠️ 动手实践
- 实现 Redis 分布式锁:在
services/agent/redis_lock.py基础上,增加可重入支持——同一 Worker 多次获取同一锁不阻塞,释放次数与获取次数匹配。 - 优化 Dockerfile:将当前 1GB 左右的 Python 镜像压缩到 200MB 以内,使用
python:3.12-slim+ 多阶段构建,并添加.dockerignore排除.git和node_modules。 - 编写 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 作用是什么?