LangGraph 生产部署指南:容器化、灰度发布与运维监控完整方案
LangGraph 生产部署指南:容器化、灰度发布与运维监控完整方案
作者简介:15年软件架构经验,曾主导多个千万级日活AI应用的从0到1落地与规模化生产,现任某头部金融科技公司AI平台技术负责人,全网LangGraph/LangChain系列技术博客累计阅读超500万。
读者定位:中级以上LangChain/LangGraph开发者、负责AI应用/Agent平台的DevOps工程师、生产环境技术架构师
阅读收获:掌握LangGraph从开发到生产全流程部署的核心技术栈、灰度发布的实用策略、全链路可观测性的落地方法,以及处理高并发、故障隔离等生产级问题的最佳实践
目录
- 核心概念与问题背景
1.1 什么是LangGraph?为什么它适合生产级Agent应用?
1.2 LangGraph开发环境与生产环境的本质差异
1.3 生产级LangGraph部署的三大核心挑战
1.4 完整生产部署方案的技术选型概览 - 容器化部署:从Dockerfile到Kubernetes
2.1 核心概念:OCI容器、容器镜像分层、多阶段构建
2.2 LangGraph应用的Dockerfile最佳实践(Python/TypeScript双语言示例)
2.3 镜像优化策略:减少体积、提升安全性、加速拉取
2.4 Kubernetes资源编排:Deployment、StatefulSet(针对有状态Agent)、Service、Ingress
2.5 配置与密钥管理:ConfigMap、Secret、外部配置中心(Apollo/Nacos/HashiCorp Vault)
2.6 依赖注入与环境隔离:LangGraph Graph实例的工厂模式实现
2.7 核心概念结构与组成要素
2.8 容器化核心要素对比表
2.9 容器化部署架构Mermaid图 - 灰度发布与流量治理:降低生产变更风险
3.1 核心概念:灰度发布、金丝雀发布、蓝绿部署、A/B测试的区别
3.2 LangGraph Agent应用的流量特点(异步流、多步骤调用、状态绑定)
3.3 流量治理的技术选型:NGINX Ingress Controller + Argo Rollouts、Istio + Knative Serving
3.4 策略驱动的灰度发布:基于用户ID/请求特征/负载均衡的规则
3.5 状态绑定Agent的灰度发布难点与解决方案(状态复制、兼容序列化、版本迁移工具)
3.6 灰度发布的全流程Mermaid流程图
3.7 灰度发布与流量治理核心ER实体关系图 - 全链路可观测性:运维监控与故障定位
4.1 核心概念:可观测性的三大支柱(Logging、Metrics、Tracing)、OpenTelemetry
4.2 LangGraph应用的可观测性痛点(多步骤调用链、状态存储监控、LLM调用性能分析)
4.3 全链路可观测性技术栈选型:Loki(日志)、Prometheus + Grafana(指标)、Jaeger/Zipkin(追踪)
4.4 OpenTelemetry与LangGraph的深度集成(自定义Span/Metric/Log处理器、状态变化追踪、LLM调用埋点)
4.5 Prometheus告警规则配置:Pod故障、LLM超时、状态存储异常、高并发指标
4.6 可观测性架构Mermaid图
4.7 可观测性核心要素对比表 - 项目实战:金融风控对话式Agent的生产部署
5.1 项目介绍:业务场景、核心功能、技术栈
5.2 开发环境搭建:LangGraph Studio、Python虚拟环境、本地依赖
5.3 系统功能设计:Agent对话流程、状态定义、工具链
5.4 系统架构设计:从开发到生产的三层架构
5.5 系统接口设计:RESTful API、WebSocket接口(异步流)
5.6 系统核心实现源代码(容器化、灰度发布、可观测性集成)
5.7 代码解读与分析
5.8 最佳实践Tips - 行业发展与未来趋势
6.1 LangGraph部署技术的演变历史(2023-2024 Q2)
6.2 未来3-5年的发展趋势:Serverless化、边缘计算部署、AI原生编排、多Agent联邦部署
6.3 面临的挑战与解决方案展望 - 本章小结
1. 核心概念与问题背景
1.1 什么是LangGraph?为什么它适合生产级Agent应用?
核心概念定义
LangGraph 是 LangChain 团队于 2023 年 10 月正式推出的构建可控、可调试、可持久化状态多智能体(Multi-Agent)和工作流应用的框架。它的核心设计思想是:以状态机(State Machine)为核心,用有向无环图(DAG)或循环图(Cyclic Graph,这是LangGraph区别于LangChain LCEL的最大特点)来描述Agent的执行流程,明确管理状态的输入、输出、更新和持久化。
适合生产级应用的五大理由
- 状态持久化与断点续传:
- LCEL 的链式调用是无状态的(Stateless),一旦中间某个工具调用或LLM请求失败,整个链就会崩溃,无法恢复;
- LangGraph 内置了多种状态存储后端(Memory、SQLite、PostgreSQL、Redis、MongoDB),可以将每一步的执行状态序列化后存储下来,支持随时暂停、恢复、甚至回滚到任意历史状态。这对于金融风控、客户服务等需要高可靠性的生产场景至关重要——比如用户正在进行贷款申请的多轮对话,网络中断后重连可以直接从断点继续,无需重新输入所有信息。
- 循环执行与条件分支的明确控制:
- LCEL 的链式调用本质上是线性的,虽然支持
RunnableBranch做条件分支,但实现循环(比如Agent需要反复调用工具直到问题解决)需要嵌套 Runnable 或自定义回调,代码可读性差,调试困难; - LangGraph 可以显式地定义循环的起始点、终止条件和过渡逻辑,比如经典的 ReAct Agent 框架,可以用
START -> Agent -> Tool -> Agent -> END这样的循环图来实现,每一步的执行状态都可以被监控和调试。
- LCEL 的链式调用本质上是线性的,虽然支持
- 可调试性与可视化:
- LangGraph 官方提供了 LangGraph Studio(本地或云端部署的可视化工具),可以实时查看 Agent 的执行流程、状态变化、中间输出、LLM 调用详情,大大降低了生产环境的调试成本;
- LangGraph 还内置了详细的回调系统,可以与 OpenTelemetry 等可观测性工具无缝集成,生成全链路的执行日志、指标和追踪。
- 可扩展性与模块化设计:
- LangGraph 的核心组件(State、Node、Edge、Graph、Compiler、Executor、Storage)都是模块化的,可以根据生产需求进行自定义替换;
- 支持将复杂的 Graph 拆分成多个子 Graph(Subgraph),每个子 Graph 可以独立开发、测试、部署,然后通过引用的方式组合起来,类似于微服务的设计思想。
- 同步/异步双模式支持与流式输出:
- 支持同步执行(适合简单的请求-响应场景)和异步执行(适合高并发、多工具并行调用的场景);
- 支持流式输出(LLM 响应流、状态变化流),可以为用户提供实时的反馈,提升用户体验,这在聊天机器人、代码生成助手等场景中非常重要。
1.2 LangGraph开发环境与生产环境的本质差异
很多开发者在本地用 LangChain/LangGraph 跑通了一个简单的 Agent 应用,就直接扔到生产环境,结果遇到了各种问题——比如性能差、稳定性低、故障定位困难等。这是因为开发环境与生产环境的目标和约束完全不同。
我们可以用下面的对比表来清晰地展示两者的差异:
| 维度 | 开发环境 | 生产环境 |
|---|---|---|
| 核心目标 | 快速迭代、验证功能、调试代码 | 高可用性、高并发性、高安全性、可观测性、低成本 |
| 用户数量 | 1-10人(开发者/测试人员) | 数千-数百万(真实用户) |
| 负载情况 | 低负载、低并发、无峰值 | 高负载、高并发、有明显的峰值(比如工作日9-12点) |
| 状态存储 | 内存存储(Memory)、本地SQLite | 分布式存储(PostgreSQL主从/集群、Redis Sentinel/Cluster) |
| LLM调用方式 | 直接调用单个模型、无限流、无重试 | 调用多个模型池、有限流、有指数退避重试、有缓存 |
| 错误处理 | 打印错误堆栈、手动重启 | 自动故障转移、降级服务、报警通知 |
| 可观测性 | 控制台输出、断点调试 | 全链路日志、指标、追踪 |
| 配置管理 | .env 文件、硬编码配置 | 外部配置中心(支持热更新)、环境变量、Secret管理 |
| 部署方式 | 本地运行、Docker单容器 | Kubernetes集群、容器编排、灰度发布 |
| 安全性 | 无认证、无授权、无加密 | 身份认证(OAuth2.0/JWT)、权限控制(RBAC/ABAC)、HTTPS/TLS加密、LLM API密钥加密存储 |
| 成本控制 | 几乎不计成本 | 严格控制计算资源、存储资源、LLM调用成本 |
1.3 生产级LangGraph部署的三大核心挑战
基于上述开发环境与生产环境的差异,生产级 LangGraph 部署主要面临以下三大核心挑战:
挑战一:状态绑定的高并发与高可用性
LangGraph 的核心优势是状态持久化,但这也带来了一个新的问题:状态绑定(State Binding)。也就是说,同一个用户的对话流程,必须由同一个 Pod 或同一个有状态服务实例来处理吗?如果是无状态的 LCEL 链式调用,我们可以用 Kubernetes Deployment 来部署,水平扩展多个 Pod,然后通过 Service 进行负载均衡;但如果是有状态的 LangGraph Agent 应用(使用内存存储或本地 SQLite 存储状态),我们就不能随便水平扩展,否则同一个用户的请求可能会被分发到不同的 Pod,导致状态丢失。
即使我们使用了分布式状态存储(比如 PostgreSQL 或 Redis),仍然面临并发更新状态的问题——比如同一个用户同时发起了两个请求,或者同一个请求的多个异步子任务同时更新状态,这可能会导致状态的不一致性(Race Condition)。
另外,状态存储本身的高可用性也是一个问题——如果 PostgreSQL 主节点宕机,我们需要能够快速切换到从节点,或者使用集群模式来保证服务的连续性。
挑战二:Agent应用的灰度发布与流量治理
传统的 Web 应用灰度发布比较简单——我们只需要根据用户ID、IP地址、请求头信息等特征,将一部分流量分发到新版本的 Pod 即可。但 LangGraph Agent 应用的灰度发布有两个特殊的难点:
- 异步流与状态绑定的流量一致性:
- 如果用户的对话流程是异步的(比如使用 WebSocket),并且状态存储在分布式数据库中,那么用户的前几个请求可能被分发到旧版本的 Pod,后面的请求又被分发到新版本的 Pod。这时候就会出现版本兼容性问题——旧版本的 Pod 写入的状态格式,新版本的 Pod 可能无法读取,或者新版本的 Pod 写入的状态格式,旧版本的 Pod 无法读取,导致整个对话流程崩溃。
- 多步骤调用链的性能与稳定性评估:
- 传统的 Web 应用通常只有 1-3 个内部调用(比如数据库查询、缓存查询、第三方 API 调用),性能与稳定性评估比较简单;但 LangGraph Agent 应用的调用链通常很长(比如 ReAct Agent 可能会反复调用 LLM 和工具 5-10 次,甚至更多),而且每次调用的性能和稳定性都可能不同(比如 LLM 调用可能会超时,工具调用可能会失败)。因此,我们需要一种更精细的灰度发布策略,能够实时监控新版本调用链的每一步的性能和稳定性,一旦出现问题,能够快速回滚。
挑战三:全链路可观测性的落地
传统的 Web 应用的可观测性已经比较成熟了——我们可以用 Nginx 访问日志、应用日志、数据库慢查询日志、Prometheus 指标(比如 QPS、响应时间、错误率)、Jaeger/Zipkin 追踪来定位问题。但 LangGraph Agent 应用的可观测性有三个特殊的痛点:
- 多步骤调用链的追踪难度大:
- 传统的 OpenTelemetry 埋点通常只能追踪到单个服务的调用,但 LangGraph Agent 应用的调用链是在同一个服务内部的多个 Node 之间跳转的,而且每个 Node 可能会调用多个外部服务(比如 LLM API、工具 API、数据库)。因此,我们需要一种更精细的追踪机制,能够将同一个对话流程的所有 Node 调用、外部服务调用都关联到同一个 Trace ID 上,并且能够展示每个 Node 的输入、输出、状态变化、执行时间等详细信息。
- 状态存储的监控不足:
- 状态存储是 LangGraph Agent 应用的核心组件,但很多开发者往往忽视了对状态存储的监控——比如状态的写入/读取延迟、错误率、存储容量、并发连接数等。如果状态存储出现问题,整个 Agent 应用都会崩溃。
- LLM调用的性能与成本分析:
- LLM 调用是 LangGraph Agent 应用的核心成本和性能瓶颈,但传统的可观测性工具通常无法提供 LLM 调用的详细信息——比如 Token 的使用量、Prompt 的长度、Completion 的长度、调用的模型名称、调用的时间、调用的成功率等。我们需要一种专门的 LLM 可观测性工具,能够实时监控 LLM 调用的性能和成本,并且能够根据监控数据优化 Prompt 设计、模型选择、缓存策略等。
1.4 完整生产部署方案的技术选型概览
为了解决上述三大核心挑战,我们设计了一套完整的 LangGraph 生产部署方案,技术栈如下:
| 技术领域 | 核心技术选型 | 备选技术选型 |
|---|---|---|
| 编程语言 | Python(主流,生态丰富) | TypeScript(适合前端全栈开发) |
| 容器化 | Docker + BuildKit(多阶段构建) | Podman + Buildah |
| 容器编排 | Kubernetes(主流,生态丰富) | Docker Swarm(小规模部署)、Nomad(HashiCorp生态) |
| 状态存储 | PostgreSQL(主从/集群,适合结构化状态) + Redis Sentinel/Cluster(缓存临时状态,提升性能) | MongoDB(适合非结构化状态)、DynamoDB(AWS生态) |
| 配置与密钥管理 | ConfigMap + Secret(K8s原生) + HashiCorp Vault(外部密钥管理) | Apollo(携程开源)、Nacos(阿里巴巴开源) |
| API网关与流量治理 | NGINX Ingress Controller(K8s原生) + Argo Rollouts(策略驱动的灰度发布) | Istio(服务网格,适合复杂的微服务架构) + Knative Serving(Serverless) |
| 全链路可观测性 | OpenTelemetry(统一埋点) + Loki(日志) + Prometheus + Grafana(指标与可视化) + Jaeger(追踪) | Datadog(商业)、New Relic(商业)、Splunk(商业) |
| 身份认证与权限控制 | OAuth2.0 + JWT + Keycloak(开源身份认证服务) | Auth0(商业)、Okta(商业) |
| LLM调用管理 | LangSmith(LangChain官方,可观测性+调试+评估) + LiteLLM(统一LLM API接口,模型池管理,缓存,限流,重试) | Langfuse(开源LLM可观测性)、Helicone(LLM API网关) |
1.5 章节核心概念与问题小结
本章节主要介绍了 LangGraph 的核心概念、适合生产级应用的原因、开发环境与生产环境的本质差异、生产级部署的三大核心挑战,以及完整生产部署方案的技术选型概览。这些内容是后续章节的基础,希望读者能够认真理解。
2. 容器化部署:从Dockerfile到Kubernetes
2.1 核心概念:OCI容器、容器镜像分层、多阶段构建
核心概念定义
- OCI容器:
- OCI(Open Container Initiative)是由 Linux 基金会主导的一个开源项目,旨在制定容器的标准规范,包括运行时规范(Runtime Specification) 和 镜像规范(Image Specification)。
- OCI 容器是一个轻量级、可移植、自包含的软件包,包含了运行一个应用所需的所有内容——代码、运行时环境、系统工具、系统库、配置文件等。
- 常见的 OCI 容器运行时包括:runc(Docker 原生的运行时)、crun(Red Hat 开源的运行时,更快更轻量)、youki(Rust 编写的运行时)。
- 容器镜像分层:
- OCI 容器镜像不是一个单一的文件,而是由多个只读的层(Layer) 组成的,每个层对应 Dockerfile 中的一条指令(比如
FROM、RUN、COPY、ADD)。 - 当我们构建一个新的容器镜像时,Docker 会复用之前已经构建好的层,只有当 Dockerfile 中的指令发生变化时,才会重新构建对应的层。这大大提高了镜像构建的速度,减少了镜像的存储空间。
- 当我们运行一个容器时,Docker 会在所有只读层的上面添加一个可写层(Writable Layer),容器运行过程中产生的所有数据(比如日志、临时文件、状态数据)都会存储在这个可写层中。当容器被删除时,这个可写层也会被删除,所有数据都会丢失——因此,我们不应该将持久化数据存储在容器的可写层中,而应该使用数据卷(Volume) 或 持久化卷声明(PersistentVolumeClaim,PVC) 来存储持久化数据。
- OCI 容器镜像不是一个单一的文件,而是由多个只读的层(Layer) 组成的,每个层对应 Dockerfile 中的一条指令(比如
- 多阶段构建(Multi-Stage Build):
- 多阶段构建是 Docker 17.05 版本引入的一个特性,允许我们在一个 Dockerfile 中使用多个
FROM指令,每个FROM指令对应一个构建阶段(Build Stage)。 - 我们可以在第一个构建阶段中使用一个包含完整编译工具链的基础镜像(比如
python:3.11-slim-bookworm或者node:20-bullseye-slim)来构建应用程序,然后在第二个构建阶段中使用一个更小的基础镜像(比如python:3.11-alpine3.19或者node:20-alpine3.19)来运行应用程序,只复制第一个构建阶段中构建好的应用程序和必要的依赖,而不复制编译工具链和其他不必要的文件。 - 多阶段构建可以大大减少容器镜像的体积——比如 Python 应用的镜像体积可以从 1GB+ 减少到 100MB-,TypeScript 应用的镜像体积可以从 2GB+ 减少到 200MB-。这不仅可以节省存储空间,还可以加速镜像的拉取和推送速度,提高部署的效率。
- 多阶段构建是 Docker 17.05 版本引入的一个特性,允许我们在一个 Dockerfile 中使用多个
2.2 LangGraph应用的Dockerfile最佳实践(Python/TypeScript双语言示例)
Python LangGraph应用的Dockerfile最佳实践
我们以一个简单的 ReAct Agent 应用 为例,来展示 Python LangGraph 应用的 Dockerfile 最佳实践。这个 ReAct Agent 可以回答用户的问题,并且可以调用 DuckDuckGoSearchRun 工具来搜索互联网上的信息。
项目结构
langgraph-react-agent/
├── .dockerignore # Docker构建时忽略的文件
├── Dockerfile # Docker镜像构建文件
├── requirements.txt # Python依赖
├── requirements-dev.txt # Python开发依赖(仅在开发阶段使用)
├── src/
│ ├── __init__.py
│ ├── main.py # 应用入口文件
│ ├── agent/
│ │ ├── __init__.py
│ │ ├── state.py # Agent状态定义
│ │ ├── nodes.py # Agent节点定义
│ │ ├── graph.py # Agent图定义
│ │ └── factory.py # Agent图工厂模式实现
│ └── config/
│ ├── __init__.py
│ └── settings.py # 配置管理(使用Pydantic Settings)
└── tests/ # 测试文件(仅在开发阶段使用)
├── __init__.py
├── test_agent.py
└── conftest.py
.dockerignore 文件
首先,我们需要编写一个 .dockerignore 文件,来指定 Docker 构建时应该忽略哪些文件,这样可以减少镜像构建的上下文(Build Context)的大小,提高镜像构建的速度,并且避免将不必要的文件(比如 .env、__pycache__、tests/、.git/)复制到镜像中,提高镜像的安全性。
# Python相关
__pycache__/
*.py[cod]
*$py.class
*.so
.Python
build/
develop-eggs/
dist/
downloads/
eggs/
.eggs/
lib/
lib64/
parts/
sdist/
var/
wheels/
*.egg-info/
.installed.cfg
*.egg
MANIFEST
# Virtual Environment
venv/
ENV/
env/
# IDE相关
.vscode/
.idea/
*.swp
*.swo
*~
# 测试相关
tests/
.pytest_cache/
.coverage
htmlcov/
# Git相关
.git/
.gitignore
.gitattributes
# 配置相关(敏感信息,不应该复制到镜像中)
.env
.env.*
!.env.example
# 其他
*.log
.DS_Store
requirements.txt 文件
接下来,我们需要编写一个 requirements.txt 文件,来指定 Python 应用的生产依赖。我们应该尽可能地使用固定版本号的依赖,而不是使用 >= 或 ~= 这样的版本范围,这样可以避免依赖的版本更新导致应用程序出现不可预测的问题。
# LangChain/LangGraph核心依赖
langchain==0.2.16
langchain-core==0.2.38
langchain-openai==0.1.23
langgraph==0.2.31
langchain-community==0.2.16 # DuckDuckGoSearchRun工具在这里
# 配置管理
pydantic-settings==2.5.2
python-dotenv==1.0.1
# Web框架(提供RESTful API和WebSocket接口)
fastapi==0.112.2
uvicorn[standard]==0.30.6
websockets==13.0
# 异步数据库驱动(PostgreSQL)
asyncpg==0.30.0
sqlalchemy==2.0.35
alembic==1.13.2
# 异步Redis驱动
redis==7.7.4
hiredis==2.3.2
# OpenTelemetry核心依赖(可观测性)
opentelemetry-api==1.27.0
opentelemetry-sdk==1.27.0
opentelemetry-instrumentation-fastapi==0.48b0
opentelemetry-instrumentation-requests==0.48b0
opentelemetry-instrumentation-sqlalchemy==0.48b0
opentelemetry-instrumentation-redis==0.48b0
opentelemetry-exporter-otlp-proto-grpc==1.27.0
opentelemetry-exporter-prometheus==0.48b0
# 日志格式化
python-json-logger==2.0.7
requirements-dev.txt 文件
然后,我们需要编写一个 requirements-dev.txt 文件,来指定 Python 应用的开发依赖(比如测试框架、代码格式化工具、代码检查工具等)。这些依赖仅在开发阶段使用,不应该复制到生产镜像中。
-r requirements.txt
# 测试框架
pytest==8.3.3
pytest-asyncio==0.24.0
pytest-cov==5.0.0
# 代码格式化工具
black==24.8.0
isort==5.13.2
# 代码检查工具
flake8==7.1.1
mypy==1.11.2
pylint==3.2.7
# LangGraph Studio
langgraph-cli==0.2.31
src/config/settings.py 文件
接下来,我们需要编写一个配置管理文件,使用 Pydantic Settings 来管理应用程序的配置。Pydantic Settings 可以从环境变量、.env 文件、外部配置中心等多个来源读取配置,并且支持配置的类型检查和验证。
from typing import Optional, List
from pydantic_settings import BaseSettings, SettingsConfigDict
from pydantic import Field, SecretStr, field_validator
class Settings(BaseSettings):
"""应用程序配置管理"""
model_config = SettingsConfigDict(
env_file=".env", # 本地开发时读取.env文件
env_file_encoding="utf-8",
case_sensitive=False, # 环境变量不区分大小写
extra="ignore", # 忽略未定义的环境变量
)
# ==================== 应用程序基础配置 ====================
app_name: str = Field(default="LangGraph ReAct Agent", description="应用程序名称")
app_version: str = Field(default="1.0.0", description="应用程序版本")
debug: bool = Field(default=False, description="是否开启调试模式")
log_level: str = Field(default="INFO", description="日志级别")
# ==================== Web服务器配置 ====================
host: str = Field(default="0.0.0.0", description="Web服务器监听地址")
port: int = Field(default=8000, description="Web服务器监听端口")
workers: int = Field(default=1, description="Uvicorn工作进程数(生产环境建议设置为CPU核心数的2-4倍,但注意LangGraph状态存储的并发问题)")
# ==================== LLM配置 ====================
openai_api_key: SecretStr = Field(..., description="OpenAI API密钥")
openai_api_base: str = Field(default="https://api.openai.com/v1", description="OpenAI API基础URL")
openai_model_name: str = Field(default="gpt-4o-mini", description="OpenAI模型名称")
openai_temperature: float = Field(default=0.0, description="OpenAI模型温度")
openai_max_tokens: Optional[int] = Field(default=None, description="OpenAI模型最大Token数")
openai_timeout: float = Field(default=30.0, description="OpenAI API调用超时时间(秒)")
openai_max_retries: int = Field(default=3, description="OpenAI API调用最大重试次数")
# ==================== 状态存储配置(PostgreSQL) ====================
postgres_host: str = Field(..., description="PostgreSQL主机地址")
postgres_port: int = Field(default=5432, description="PostgreSQL端口")
postgres_user: str = Field(..., description="PostgreSQL用户名")
postgres_password: SecretStr = Field(..., description="PostgreSQL密码")
postgres_db: str = Field(..., description="PostgreSQL数据库名称")
postgres_pool_size: int = Field(default=10, description="PostgreSQL连接池大小")
postgres_max_overflow: int = Field(default=20, description="PostgreSQL连接池最大溢出数")
postgres_pool_timeout: float = Field(default=30.0, description="PostgreSQL连接池超时时间(秒)")
postgres_pool_recycle: float = Field(default=3600.0, description="PostgreSQL连接池回收时间(秒)")
@field_validator("postgres_db")
@classmethod
def validate_postgres_db(cls, v: str) -> str:
"""验证PostgreSQL数据库名称"""
if not v:
raise ValueError("PostgreSQL数据库名称不能为空")
return v
# ==================== 临时状态缓存配置(Redis) ====================
redis_host: str = Field(..., description="Redis主机地址")
redis_port: int = Field(default=6379, description="Redis端口")
redis_password: Optional[SecretStr] = Field(default=None, description="Redis密码")
redis_db: int = Field(default=0, description="Redis数据库编号")
redis_pool_size: int = Field(default=10, description="Redis连接池大小")
redis_max_connections: Optional[int] = Field(default=20, description="Redis最大连接数")
redis_socket_timeout: float = Field(default=5.0, description="Redis socket超时时间(秒)")
redis_socket_connect_timeout: float = Field(default=5.0, description="Redis socket连接超时时间(秒)")
redis_cache_ttl: int = Field(default=3600, description="Redis临时状态缓存TTL(秒)")
# ==================== OpenTelemetry配置 ====================
otel_enabled: bool = Field(default=True, description="是否开启OpenTelemetry可观测性")
otel_service_name: str = Field(default="langgraph-react-agent", description="OpenTelemetry服务名称")
otel_service_version: str = Field(default="1.0.0", description="OpenTelemetry服务版本")
otel_exporter_otlp_endpoint: str = Field(default="http://otel-collector:4317", description="OpenTelemetry OTLP导出器端点(gRPC)")
otel_exporter_prometheus_port: int = Field(default=9464, description="OpenTelemetry Prometheus导出器端口")
# ==================== 配置中心配置(可选,HashiCorp Vault) ====================
vault_enabled: bool = Field(default=False, description="是否开启HashiCorp Vault配置中心")
vault_url: str = Field(default="http://vault:8200", description="HashiCorp Vault地址")
vault_token: Optional[SecretStr] = Field(default=None, description="HashiCorp Vault令牌")
vault_secret_path: str = Field(default="secret/langgraph-react-agent", description="HashiCorp Vault密钥路径")
# 全局配置实例
settings = Settings()
src/agent/state.py 文件
然后,我们需要定义 Agent 的状态。LangGraph 的状态可以是任意的 Python 对象,但推荐使用 Pydantic BaseModel 或 TypedDict 来定义,这样可以获得更好的类型检查和验证。
from typing import List, TypedDict, Annotated
from langchain_core.messages import BaseMessage, HumanMessage, AIMessage, ToolMessage
from langgraph.graph.message import add_messages
# 使用TypedDict定义Agent状态,Annotated用于指定状态更新的Reducer函数
class AgentState(TypedDict):
"""Agent状态定义"""
# messages: 对话历史,使用add_messages作为Reducer函数,自动合并新的消息
messages: Annotated[List[BaseMessage], add_messages]
# user_id: 用户ID,用于绑定状态
user_id: str
# session_id: 会话ID,用于绑定状态
session_id: str
# iteration_count: 迭代次数,用于限制Agent的循环次数
iteration_count: int
# max_iterations: 最大迭代次数,用于限制Agent的循环次数
max_iterations: int
# tool_responses: 工具响应列表,用于存储工具调用的结果(可选,也可以直接存储在messages中)
tool_responses: List[ToolMessage]
src/agent/nodes.py 文件
接下来,我们需要定义 Agent 的节点(Node)。节点是 LangGraph 图的基本执行单元,每个节点对应一个 Python 函数,接收状态作为输入,返回一个字典作为输出(这个字典会被用来更新状态)。
from typing import Dict, Any, Literal
from langchain_core.messages import BaseMessage, HumanMessage, AIMessage, ToolMessage, SystemMessage
from langchain_core.tools import BaseTool, tool
from langchain_openai import ChatOpenAI
from langgraph.prebuilt import ToolNode
from langchain_community.tools import DuckDuckGoSearchRun
from src.config.settings import settings
from src.agent.state import AgentState
# ==================== 工具定义 ====================
@tool
def duckduckgo_search(query: str) -> str:
"""使用DuckDuckGo搜索互联网上的信息"""
search = DuckDuckGoSearchRun()
return search.run(query)
# 工具列表
tools: List[BaseTool] = [duckduckgo_search]
# 工具节点(LangGraph预构建的节点,用于调用工具)
tool_node = ToolNode(tools)
# ==================== LLM初始化 ====================
def get_llm() -> ChatOpenAI:
"""初始化LLM"""
return ChatOpenAI(
api_key=settings.openai_api_key.get_secret_value(),
base_url=settings.openai_api_base,
model=settings.openai_model_name,
temperature=settings.openai_temperature,
max_tokens=settings.openai_max_tokens,
timeout=settings.openai_timeout,
max_retries=settings.openai_max_retries,
).bind_tools(tools) # 绑定工具
# ==================== 节点定义 ====================
def agent_node(state: AgentState) -> Dict[str, Any]:
"""
Agent节点:调用LLM生成响应或工具调用请求
"""
# 初始化LLM
llm = get_llm()
# 系统提示词
system_prompt = SystemMessage(content="""你是一个 helpful 的 AI 助手,可以回答用户的问题。
如果你不知道答案,可以使用 duckduckgo_search 工具来搜索互联网上的信息。
请尽可能简洁、准确地回答用户的问题。""")
# 构建消息列表:系统提示词 + 对话历史
messages = [system_prompt] + state["messages"]
# 调用LLM
response: AIMessage = llm.invoke(messages)
# 更新迭代次数
new_iteration_count = state["iteration_count"] + 1
# 返回状态更新字典
return {
"messages": [response],
"iteration_count": new_iteration_count,
}
def should_continue(state: AgentState) -> Literal["tools", "__end__"]:
"""
条件边(Conditional Edge):判断Agent是否应该继续调用工具或结束对话
"""
# 检查迭代次数是否超过最大限制
if state["iteration_count"] >= state["max_iterations"]:
return "__end__"
# 检查最后一条消息是否是工具调用请求
last_message: BaseMessage = state["messages"][-1]
if isinstance(last_message, AIMessage) and last_message.tool_calls:
return "tools"
# 否则结束对话
return "__end__"
src/agent/graph.py 文件
然后,我们需要定义 Agent 的图(Graph)。我们使用 StateGraph 来定义有状态的图,使用 START 和 END 节点来表示图的起始和结束,使用 add_node 方法来添加节点,使用 add_edge 方法来添加普通边,使用 add_conditional_edges 方法来添加条件边。
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver
from langgraph.checkpoint.aiosqlite import AsyncSqliteSaver
from sqlalchemy.ext.asyncio import create_async_engine, AsyncEngine
from src.config.settings import settings
from src.agent.state import AgentState
from src.agent.nodes import agent_node, tool_node, should_continue
def create_graph(checkpointer=None) -> StateGraph:
"""
创建Agent图
:param checkpointer: 状态检查点(Checkpoint)持久化对象
:return: 编译后的Agent图
"""
# 初始化StateGraph
workflow = StateGraph(AgentState)
# 添加节点
workflow.add_node("agent", agent_node)
workflow.add_node("tools", tool_node)
# 添加边
workflow.add_edge(START, "agent")
workflow.add_conditional_edges(
"agent",
should_continue,
{
"tools": "tools",
"__end__": END,
},
)
workflow.add_edge("tools", "agent")
# 编译图,添加checkpointer(如果有的话)
graph = workflow.compile(checkpointer=checkpointer)
return graph
# ==================== 异步PostgreSQL Checkpointer初始化 ====================
async def get_postgres_checkpointer() -> AsyncPostgresSaver:
"""
初始化异步PostgreSQL Checkpointer
:return: AsyncPostgresSaver对象
"""
# 构建PostgreSQL连接字符串
postgres_url = (
f"postgresql+asyncpg://{settings.postgres_user}:{settings.postgres_password.get_secret_value()}"
f"@{settings.postgres_host}:{settings.postgres_port}/{settings.postgres_db}"
)
# 创建异步SQLAlchemy引擎
engine: AsyncEngine = create_async_engine(
postgres_url,
pool_size=settings.postgres_pool_size,
max_overflow=settings.postgres_max_overflow,
pool_timeout=settings.postgres_pool_timeout,
pool_recycle=settings.postgres_pool_recycle,
echo=settings.debug,
)
# 初始化AsyncPostgresSaver
checkpointer = AsyncPostgresSaver(engine)
# 创建Checkpoint表(如果不存在的话)
async with engine.begin() as conn:
await checkpointer.setup(conn)
return checkpointer
# ==================== 异步SQLite Checkpointer初始化(仅用于本地开发) ====================
async def get_sqlite_checkpointer() -> AsyncSqliteSaver:
"""
初始化异步SQLite Checkpointer(仅用于本地开发)
:return: AsyncSqliteSaver对象
"""
checkpointer = AsyncSqliteSaver.from_conn_string("checkpoints.sqlite")
await checkpointer.setup()
return checkpointer
src/agent/factory.py 文件
接下来,我们需要实现 Agent 图的工厂模式。工厂模式可以帮助我们根据不同的环境(开发、测试、生产)初始化不同的 Checkpointer,并且可以支持依赖注入和热更新。
from typing import Optional
from langgraph.graph import StateGraph
from src.config.settings import settings
from src.agent.graph import create_graph, get_postgres_checkpointer, get_sqlite_checkpointer
class AgentGraphFactory:
"""Agent图工厂类"""
_instance: Optional["AgentGraphFactory"] = None
_graph: Optional[StateGraph] = None
def __new__(cls) -> "AgentGraphFactory":
"""单例模式:确保整个应用程序只有一个AgentGraphFactory实例"""
if cls._instance is None:
cls._instance = super().__new__(cls)
return cls._instance
async def initialize_graph(self) -> None:
"""初始化Agent图"""
if self._graph is not None:
return
# 根据环境选择Checkpointer
if settings.debug:
# 本地开发环境:使用SQLite Checkpointer
checkpointer = await get_sqlite_checkpointer()
else:
# 生产环境:使用PostgreSQL Checkpointer
checkpointer = await get_postgres_checkpointer()
# 创建Agent图
self._graph = create_graph(checkpointer=checkpointer)
def get_graph(self) -> StateGraph:
"""获取Agent图"""
if self._graph is None:
raise RuntimeError("Agent图尚未初始化,请先调用initialize_graph()方法")
return self._graph
async def reset_graph(self) -> None:
"""重置Agent图(用于热更新)"""
self._graph = None
await self.initialize_graph()
# 全局Agent图工厂实例
agent_graph_factory = AgentGraphFactory()
src/main.py 文件
然后,我们需要编写应用程序的入口文件,使用 FastAPI 来提供 RESTful API 和 WebSocket 接口。
import logging
from typing import Dict, Any
from contextlib import asynccontextmanager
from fastapi import FastAPI, Request, WebSocket, WebSocketDisconnect, HTTPException, status
from fastapi.responses import JSONResponse
from langchain_core.messages import HumanMessage
from src.config.settings import settings
from src.agent.state import AgentState
from src.agent.factory import agent_graph_factory
from src.utils.logger import setup_logger
from src.utils.otel import setup_otel
# ==================== 日志初始化 ====================
logger = setup_logger(settings.app_name, settings.log_level)
# ==================== OpenTelemetry初始化 ====================
if settings.otel_enabled:
setup_otel(
service_name=settings.otel_service_name,
service_version=settings.otel_service_version,
otlp_endpoint=settings.otel_exporter_otlp_endpoint,
prometheus_port=settings.otel_exporter_prometheus_port,
)
# ==================== FastAPI生命周期管理 ====================
@asynccontextmanager
async def lifespan(app: FastAPI):
"""FastAPI生命周期管理:应用启动时初始化Agent图,应用关闭时清理资源"""
logger.info(f"正在启动应用程序 {settings.app_name} v{settings.app_version}...")
# 初始化Agent图
try:
await agent_graph_factory.initialize_graph()
logger.info("Agent图初始化成功")
except Exception as e:
logger.error(f"Agent图初始化失败: {str(e)}", exc_info=True)
raise e
yield # 应用程序运行期间
# 清理资源
logger.info(f"正在关闭应用程序 {settings.app_name} v{settings.app_version}...")
# ==================== FastAPI应用初始化 ====================
app = FastAPI(
title=settings.app_name,
version=settings.app_version,
debug=settings.debug,
lifespan=lifespan,
)
# ==================== 全局异常处理 ====================
@app.exception_handler(Exception)
async def global_exception_handler(request: Request, exc: Exception):
"""全局异常处理:返回统一的错误响应格式"""
logger.error(f"未捕获的异常: {str(exc)}", exc_info=True)
return JSONResponse(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
content={
"status": "error",
"message": "内部服务器错误",
"detail": str(exc) if settings.debug else None,
},
)
# ==================== 健康检查接口 ====================
@app.get("/health", status_code=status.HTTP_200_OK)
async def health_check() -> Dict[str, Any]:
"""健康检查接口:用于Kubernetes的Liveness Probe和Readiness Probe"""
return {
"status": "healthy",
"app_name": settings.app_name,
"app_version": settings.app_version,
}
# ==================== RESTful API接口 =================
更多推荐


所有评论(0)