LangGraph多智能体系统运维:从部署到监控的自动化方案

你好,我是Alex,一位专注于LLMOps领域深耕3年、踩过无数LangChain/LangGraph部署坑的全栈工程师。今天这篇文章,我会把从0搭建LangGraph企业级集群自动化部署平台、再搭建端到端多维度监控体系的完整经验,毫无保留地分享给你。


引言

痛点引入

上个月帮我前同事老张,哦对了他是国内某跨境电商公司的AI中台负责人,凌晨三点在咱们LLMOps群里炸锅了:

“救命救命,LangGraph做的那个选品+投诉自动回复+供应链调度的**三元协同智能体,供应链节点连续两天在深夜崩溃重启4次!重启前完全没预警,重启过程还要手动SSH进12台容器挨个改配置同步状态??研发赶周报时我差点跑路!

你上次说要做的那个全链路自动化运维方案,能不能发我?

是啊,LangGraph虽然让多智能体应用从“原型跑通Demo”到“企业级生产环境可用”之间,隔了整整一个从0到1搭建自动化平台的鸿沟:

  1. **部署痛点:
    • 原型代码依赖复杂:Python环境版本死锁、LangChain/LangGraph/Vector DB(Pinecone/Chroma本地→本地→Milvus/Zilliz云原生的依赖链混乱;
    • 单机部署卡资源:单智能体调度节点、状态持久化Redis/Milvus/Zilliz、向量嵌入节点、工具调用节点挤在一台8核16G的ECS上,高峰期并发100QPS就CPU飙升到99%;
    • 集群部署难管理:手动启动10台容器同步状态Redis哨兵、配置LangGraph分布式调度器、监控节点健康度要写一堆Shell脚本,写一次用一次改配置又要重写;
    • 灰度发布不敢试:直接全量换版本,三元协同智能体哪怕一个工具调用Agent的Prompt优化出问题,整个选品链路全挂;
  2. **监控痛点:
    • 智能体行为不可见:只知道调用总QPS和错误率,不知道三元协同智能体里的每个Agent谁先超时、向量检索花了多久、状态持久化Redis的键值冲突多少次;
    • 告警不及时:容器挂了10分钟才收到Kubernetes Pod CrashLoopBackOff的告警,但此时投诉处理已经超时1000条;
    • 性能瓶颈难定位:高峰期响应慢,不知道是GPU资源不够嵌入向量、还是Redis哨兵同步慢、还是供应链调度工具的第三方API超时;
    • 成本无法控制:不知道哪些Agent调用工具最多、哪些嵌入节点闲置,盲目加资源每月成本翻了3倍。

解决方案概述

针对这些痛点,我和我的团队在LangGraph 0.2.4(截至202X年Q2的最新稳定版)基础上,设计并落地了一套企业级LangGraph多智能体系统全链路自动化运维方案

  1. **部署层:
    • 使用Docker Compose快速搭建本地开发/测试环境;
    • 使用Helm Charts一键部署Kubernetes(K8s)生产环境集群,解决资源隔离、负载均衡、状态同步;
    • 基于Argo CD实现GitOps自动化部署和金丝雀发布;
    • 基于OpenTofu实现基础设施即代码(IaC)自动化管理云资源;
  2. **监控层:
    • 基于Prometheus + Grafana实现资源监控和性能监控;
    • 基于LangSmith(LangChain官方) + OpenTelemetry实现可观测性(Traces、Metrics、Logs);
    • 基于Alertmanager实现多渠道告警;
    • 基于LlamaIndex Observability + LangSmith集成实现更细粒度的工具调用、向量检索、Agent决策监控;
  3. **自动化层:
    • 基于Kubernetes Horizontal Pod Autoscaler(HPA)实现节点自动扩缩容;
    • 基于Argo Workflows实现自动化故障排查和自愈;
    • 基于CI/CD Pipeline实现从代码提交到生产环境全自动化部署测试。

最终效果展示

这套方案在我前同事老张的跨境电商AI中台上线后:

  1. **部署效率提升:
    • 从手动部署12台容器到一键部署100台容器,时间从3小时缩短到5分钟;
    • 金丝雀发布成功率从60%提升到98%;
  2. **监控响应提升:
    • 端到端可观测性覆盖率从10%提升到100%;
    • 告警提前预警率从10%提升到95%;
  3. **可用性提升:
    • 三元协同智能体的可用性从92%提升到99.9%;
  4. **成本降低:
    • 每月云资源成本降低了40%。

(为了让大家更直观地看到效果,我在GitHub上写了一个简化版的三元协同选品智能体,配套完整的Docker Compose、Helm Charts、Prometheus+Grafana配置、LangSmith集成、OpenTelemetry集成,只需要10分钟就能跑通本地开发环境和监控体系的全链路。项目地址:https://github.com/alex-llmops/langgraph-ops-demo


准备工作

环境/工具

在开始搭建这套方案之前,你需要准备好以下开发环境和工具:

本地开发环境
工具/依赖 版本要求 说明
操作系统 macOS 12+/Ubuntu 22.04+/Windows 11+ WSL2 推荐使用macOS或Ubuntu,Windows建议使用WSL2避免环境问题
Docker 20.10+ 用于构建和运行Docker容器
Docker Compose 2.20+ 用于快速搭建本地开发/测试环境
Python 3.10-3.12 LangGraph 0.2.4稳定版仅支持Python 3.10-3.12
pip 23.0+ Python包管理器,推荐使用uv(更高效的Python包管理器)
uv 0.2.0+ 更高效的Python包管理器,用于快速安装依赖和管理虚拟环境
Git 2.40+ 代码版本控制
VS Code 1.85+ 代码编辑器,推荐安装以下插件:
→ Python插件 最新版 支持Python开发
→ Docker插件 最新版 支持Docker开发
→ Kubernetes插件 最新版 支持Kubernetes开发
→ LangChain插件 最新版 支持LangChain/LangGraph开发
生产环境(可选,但强烈推荐)
工具/依赖 版本要求 说明
Kubernetes集群 1.28-1.30 可以使用EKS(AWS)、GKE(GCP)、AKS(Azure)、Rancher或minikube(测试用)
Helm 3.13+ Kubernetes包管理器,用于一键部署应用和服务
Argo CD 2.10+ GitOps工具,用于自动化部署和金丝雀发布
OpenTofu/Terraform 1.6+ 基础设施即代码工具,用于自动化管理云资源
Prometheus 2.45+ 时序数据库,用于存储和查询监控指标
Grafana 10.4+ 可视化工具,用于展示监控仪表盘
Alertmanager 0.26+ 告警管理工具,用于多渠道告警
LangSmith API Key 免费版/专业版 LangChain官方提供的可观测性平台,用于监控智能体行为
OpenTelemetry Collector 0.90+ 可观测性数据收集和转发工具,用于将数据转发到Prometheus、Grafana Loki、Jaeger或LangSmith
Grafana Loki 2.9+ 日志聚合工具,用于存储和查询日志
Jaeger 1.53+ 分布式追踪工具,用于追踪请求链路
Milvus/Zilliz Cloud 2.3+ 向量数据库,用于存储和查询向量嵌入
Redis Sentinel/RediSearch 7.2+ 状态持久化数据库,用于存储LangGraph的状态

基础知识

为了更好地理解这套方案,你需要具备以下前置知识:

  1. **Python基础:掌握Python的基本语法、函数、类、装饰器、异步编程等;
  2. **LangChain/LangGraph基础:掌握LangChain的LLM、Prompt Template、Chain、Tool、Agent的基本概念,掌握LangGraph的State、Node、Edge、Compiled Graph的基本概念;
  3. **Docker基础:掌握Docker的基本命令、Dockerfile的编写、Docker Compose的使用;
  4. **Kubernetes基础:掌握Kubernetes的Pod、Deployment、Service、Ingress、ConfigMap、Secret、HPA、StatefulSet的基本概念;
  5. **GitOps基础:掌握GitOps的基本概念、Argo CD的使用;
  6. **Prometheus+Grafana基础:掌握PromQL的基本语法、Grafana仪表盘的配置;
  7. **OpenTelemetry基础:掌握OpenTelemetry的Traces、Metrics、Logs的基本概念、OpenTelemetry Collector的配置。

如果你不具备以上前置知识,我推荐你先学习以下资源:

  1. Python基础: Python官方教程
  2. LangChain/LangGraph基础: LangChain官方文档LangGraph官方文档我的LangGraph入门教程
  3. Docker基础: Docker官方教程
  4. Kubernetes基础: Kubernetes官方教程我的Kubernetes入门教程
  5. GitOps基础: Argo CD官方教程
  6. Prometheus+Grafana基础: Prometheus官方教程Grafana官方教程
  7. OpenTelemetry基础: OpenTelemetry官方教程

核心概念

问题背景

在讲LangGraph多智能体系统运维的核心概念之前,我们先得搞清楚一个问题:**为什么LangGraph多智能体系统的运维比传统的单体应用/微服务应用难?

传统的单体应用运维只需要监控服务器CPU、内存、磁盘、网络、日志、错误率就够了;传统的微服务应用运维虽然多了分布式追踪,但服务之间的调用关系是固定的(通过HTTP/gRPC定义的),状态持久化也是通过数据库/Redis固定的。

但LangGraph多智能体系统不一样:

  1. **动态调用关系:Agent之间的调用关系是由State驱动的,不是固定的HTTP/gRPC调用,同一个Request可能走不同的路径,比如投诉自动回复Agent可能先调用向量检索Agent获取相似投诉,也可能先调用工具查询Agent获取用户订单信息,也可能先调用LLM生成回复直接返回;
  2. **复杂状态管理:LangGraph的State是分布式的,可能存储在内存、Redis、PostgreSQL、MongoDB里,状态的键值冲突率、同步延迟、大小变化都会影响系统的性能和可用性;
  3. **第三方依赖不确定性:Agent可能调用大量的第三方工具(比如向量数据库、第三方API、搜索引擎、数据库、爬虫等),第三方工具的延迟、错误率、可用性都是不确定的;
  4. **LLM输出不确定性:LLM的输出是不确定的,可能超时、可能幻觉、可能不符合要求,都会影响系统的性能和可用性;
  5. **资源消耗不均衡:不同的Agent消耗的资源不一样,比如向量检索Agent消耗的GPU资源多,供应链调度Agent消耗的CPU资源多,状态持久化Redis消耗的内存资源多。

核心概念定义

为了解决这些问题,我们需要先定义LangGraph多智能体系统运维的核心概念:

1. LangGraph多智能体系统全链路可观测性(Three Pillars of Observability for LangGraph)

LangGraph多智能体系统的全链路可观测性是指能够从**Traces(分布式追踪)、Metrics(指标)、Logs(日志)**三个维度全面了解系统的运行状态,包括:

  • Traces维度: 能够追踪同一个Request在整个LangGraph多智能体系统中的完整路径,包括:
    • 每个Node的执行时间、输入输出、状态变化;
    • 每个Tool的调用时间、输入输出、错误率;
    • 每个LLM的调用时间、输入输出、Token使用量、费用、幻觉率;
    • 每个Vector DB的调用时间、输入输出、检索精度、召回率;
  • Metrics维度: 能够统计系统的性能指标、资源指标、业务指标,包括:
    • 性能指标:总QPS、总错误率、总响应时间、每个Node的QPS/错误率/响应时间、每个Tool的QPS/错误率/响应时间、每个LLM的QPS/错误率/响应时间/Token使用量/费用、每个Vector DB的QPS/错误率/响应时间/检索精度/召回率;
    • 资源指标:每个Pod/容器/节点的CPU、内存、磁盘、网络使用量,每个LLM的GPU使用量,每个Vector DB的向量索引大小;
    • 业务指标:选品智能体的选品成功率、投诉自动回复Agent的回复率、用户满意度、供应链调度Agent的调度成功率;
  • Logs维度: 能够收集系统的所有日志,包括:
    • LangGraph的运行日志;
    • LangChain的运行日志;
    • 每个Node的运行日志;
    • 每个Tool的运行日志;
    • 每个LLM的运行日志;
    • 每个Vector DB的运行日志;
    • 操作系统的运行日志;
    • Kubernetes的运行日志。
2. LangGraph多智能体系统自动化部署(Automated Deployment for LangGraph)

LangGraph多智能体系统的自动化部署是指能够从代码提交→单元测试→集成测试→预发布环境部署→生产环境金丝雀发布→全量发布的全流程自动化,包括:

  • CI/CD Pipeline: 能够自动化代码提交后的单元测试、集成测试、Docker镜像构建、Docker镜像推送;
  • GitOps: 能够自动化将预发布环境/生产环境的部署配置同步到Git仓库,然后Argo CD自动检测Git仓库的变化并自动部署;
  • 金丝雀发布: 能够先将新版本部署到一小部分Pod上,然后根据监控指标决定是否全量发布;
  • 回滚: 能够在新版本出现问题时自动回滚到旧版本。
3. LangGraph多智能体系统自动化扩缩容(Automated Scaling for LangGraph)

LangGraph多智能体系统的自动化扩缩容是指能够根据监控指标自动增加或减少Pod/容器/节点的数量,包括:

  • 水平扩缩容(HPA): 能够根据Pod的CPU/内存使用量、QPS、响应时间自动增加或减少Pod的数量;
  • 垂直扩缩容(VPA): 能够根据Pod的CPU/内存使用量自动调整Pod的CPU/内存请求和限制;
  • 节点自动扩缩容(Cluster Autoscaler): 能够根据Pod的Pending状态自动增加或减少Kubernetes集群的节点数量。
4. LangGraph多智能体系统自动化故障排查和自愈(Automated Troubleshooting and Self-Healing for LangGraph)

LangGraph多智能体系统的自动化故障排查和自愈是指能够在系统出现问题时自动排查问题自动修复问题,包括:

  • 自动化故障排查: 能够在系统出现告警时自动收集Traces、Metrics、Logs,然后自动分析问题的原因;
  • 自动化自愈: 能够在系统出现问题时自动重启Pod/容器/节点、自动回滚到旧版本、自动切换到备用节点、自动清理Redis哨兵同步慢的节点。

概念结构与核心要素组成

为了更直观地理解LangGraph多智能体系统运维的核心概念,我们可以用以下的概念结构与核心要素组成图来表示:

LangGraph多智能体系统运维全链路

核心系统层

自动化层

监控层

部署层

Docker容器

Helm Charts

Argo CD GitOps

OpenTofu IaC

CI/CD Pipeline

金丝雀发布

回滚

LangSmith 可观测性

OpenTelemetry 数据收集

Prometheus 时序数据库

Grafana 可视化

Alertmanager 告警管理

Grafana Loki 日志聚合

Jaeger 分布式追踪

LlamaIndex Observability

Horizontal Pod Autoscaler 水平扩缩容

Vertical Pod Autoscaler 垂直扩缩容

Cluster Autoscaler 节点自动扩缩容

Argo Workflows 自动化故障排查和自愈

LangGraph 多智能体系统

LLM 大语言模型

Vector DB 向量数据库

Redis Sentinel 状态持久化

第三方工具

Kubernetes

概念之间的关系

为了更深入地理解LangGraph多智能体系统运维的核心概念之间的关系,我们可以用以下的概念核心属性维度对比ER实体关系图交互关系图来表示:

概念核心属性维度对比
核心概念 核心目标 核心工具 核心指标 核心优势
全链路可观测性 全面了解系统的运行状态 LangSmith、OpenTelemetry、Prometheus、Grafana、Loki、Jaeger Traces覆盖率、Metrics覆盖率、Logs覆盖率、告警提前预警率、故障排查时间 能够快速定位问题的原因
自动化部署 快速、安全地部署新版本 Docker、Helm、Argo CD、OpenTofu、CI/CD Pipeline 部署时间、金丝雀发布成功率、回滚时间 部署效率提升、风险降低
自动化扩缩容 根据监控指标自动调整资源 HPA、VPA、Cluster Autoscaler Pod数量变化时间、资源利用率、成本降低率 资源利用率提升、成本降低
自动化故障排查和自愈 自动排查问题并自动修复 Argo Workflows、Alertmanager、Argo CD 故障自愈率、故障修复时间、系统可用性 系统可用性提升、故障修复时间缩短
ER实体关系图

包含

使用

使用

触发

收集

使用

使用

展示

使用

使用

使用

包含

存储

连接

调用

调用

调用

监控

监控

监控

监控

监控

获取

查询

查询

触发

触发

操作

操作

调整

DEPLOYMENT

DOCKER_CONTAINER

HELM_CHART

ARGOCD_APP

CICD_PIPELINE

MONITORING

OPENTELEMETRY_COLLECTOR

LANGSMITH_PROJECT

PROMETHEUS_SERVER

GRAFANA_DASHBOARD

ALERTMANAGER_RULE

AUTOMATION

HPA_OBJECT

ARGOWORKFLOW_TEMPLATE

LANGGRAPH_SYSTEM

LANGGRAPH_NODE

LANGGRAPH_STATE

LANGGRAPH_EDGE

LLM_CALL

VECTORDB_CALL

TOOL_CALL

交互关系图
LangSmith Jaeger Loki Argo Workflows Alertmanager Grafana Prometheus OpenTelemetry Collector 第三方工具 Redis Sentinel Vector DB向量数据库 LLM大语言模型 LangGraph多智能体系统 Kubernetes集群 Argo CD Git仓库 CI/CD Pipeline 用户 LangSmith Jaeger Loki Argo Workflows Alertmanager Grafana Prometheus OpenTelemetry Collector 第三方工具 Redis Sentinel Vector DB向量数据库 LLM大语言模型 LangGraph多智能体系统 Kubernetes集群 Argo CD Git仓库 CI/CD Pipeline 用户 提交代码 1 单元测试、集成测试 2 构建Docker镜像 3 推送Docker镜像到Docker Hub 4 更新Helm Charts配置 5 检测Git仓库的变化 6 部署/更新Helm Charts 7 启动/更新Pod 8 发送请求 9 读取/写入状态 10 向量检索 11 调用LLM 12 调用第三方工具 13 返回响应 14 收集Traces、Metrics、Logs 15 收集Traces、Metrics、Logs 16 收集Traces、Metrics、Logs 17 收集Traces、Metrics、Logs 18 收集Traces、Metrics、Logs 19 转发Metrics 20 转发Logs(可选,图中未画) 21 转发Traces(可选,图中未画) 22 转发Traces、Metrics、Logs 23 查询Metrics 24 查询Traces、Metrics、Logs 25 展示监控仪表盘 26 触发告警 27 发送告警通知 28 触发Argo Workflows 29 自动化故障排查 30 自动回滚 31 自动重启Pod 32 自动调整HPA 33

简化版三元协同选品智能体项目介绍

在开始搭建部署和监控体系之前,我们先得有一个简化版的三元协同选品智能体作为我们的测试项目。这个项目包含三个Agent:

  1. 选品Agent(Product Selection Agent): 根据用户的需求(比如“性价比高的无线蓝牙耳机”),调用向量检索Agent获取相似的历史选品数据,然后调用LLM生成选品方案;
  2. 价格对比Agent(Price Comparison Agent): 根据选品Agent生成的选品方案,调用第三方API(比如淘宝API、京东API、拼多多API)获取当前的价格,然后生成价格对比报告;
  3. 库存检查Agent(Inventory Check Agent): 根据价格对比Agent生成的价格对比报告,调用第三方API(比如供应商API)获取当前的库存,然后生成最终的选品推荐。

项目架构设计

简化版三元协同选品智能体的项目架构设计如下图所示:

LLM层

数据层

工具层

LangGraph层

客户端层

用户

Web UI(可选)

选品Agent

价格对比Agent

库存检查Agent

路由Agent

向量检索Tool

淘宝API Tool

京东API Tool

拼多多API Tool

供应商API Tool

Milvus向量数据库

Redis Sentinel

OpenAI GPT-4o

Anthropic Claude 3 Opus(可选)

系统接口设计

简化版三元协同选品智能体的系统接口设计如下图所示:

1. 用户请求接口(POST /api/recommendation)
  • 接口说明: 用户发送选品需求,系统返回最终的选品推荐;
  • 请求参数:
    {
      "user_id": "string",
      "request_id": "string",
      "query": "string"
    }
    
    参数名 类型 必填 说明
    user_id string 用户ID
    request_id string 请求ID,用于分布式追踪
    query string 选品需求
  • 响应参数:
    {
      "request_id": "string",
      "status": "success/failed",
      "data": {
        "product_selection_scheme": "string",
        "price_comparison_report": "string",
        "inventory_check_report": "string",
        "final_recommendation": "string"
      },
      "error": {
        "code": "string",
        "message": "string"
      }
    }
    
    参数名 类型 说明
    request_id string 请求ID,用于分布式追踪
    status string 状态:success/failed
    data.product_selection_scheme string 选品方案
    data.price_comparison_report string 价格对比报告
    data.inventory_check_report string 库存检查报告
    data.final_recommendation string 最终选品推荐
    error.code string 错误码
    error.message string 错误消息
2. 健康检查接口(GET /api/health)
  • 接口说明: 检查系统的健康状态;
  • 响应参数:
    {
      "status": "healthy/unhealthy",
      "components": {
        "langgraph": "healthy/unhealthy",
        "openai": "healthy/unhealthy",
        "milvus": "healthy/unhealthy",
        "redis": "healthy/unhealthy",
        "taobao_api": "healthy/unhealthy",
        "jd_api": "healthy/unhealthy",
        "pinduoduo_api": "healthy/unhealthy",
        "supplier_api": "healthy/unhealthy"
      }
    }
    

系统核心实现源代码

为了方便大家理解,我把简化版三元协同选品智能体的核心实现源代码放在GitHub上:https://github.com/alex-llmops/langgraph-ops-demo

下面我会把核心实现源代码的关键部分贴出来:

1. 状态定义(src/state.py)
from typing import Annotated, TypedDict, List, Optional
from langgraph.graph.message import add_messages
from langchain_core.messages import BaseMessage

class ProductRecommendationState(TypedDict):
    # 用户输入
    user_id: str
    request_id: str
    query: str
    # 消息历史
    messages: Annotated[List[BaseMessage], add_messages]
    # 选品Agent输出
    product_selection_scheme: Optional[str] = None
    # 价格对比Agent输出
    price_comparison_report: Optional[str] = None
    # 库存检查Agent输出
    inventory_check_report: Optional[str] = None
    # 最终推荐
    final_recommendation: Optional[str] = None
    # 错误信息
    error: Optional[dict] = None
2. 工具定义(src/tools.py)
from langchain_core.tools import tool
from langchain_openai import OpenAIEmbeddings
from pymilvus import MilvusClient, connections, utility
from typing import List, Dict, Optional
import os
import time
import random

# 加载环境变量
from dotenv import load_dotenv
load_dotenv()

# 初始化Milvus向量数据库连接
MILVUS_URI = os.getenv("MILVUS_URI", "http://localhost:19530")
MILVUS_COLLECTION_NAME = os.getenv("MILVUS_COLLECTION_NAME", "product_recommendations")
embeddings = OpenAIEmbeddings(model="text-embedding-3-small")
milvus_client = MilvusClient(uri=MILVUS_URI)

# 初始化向量数据库集合(如果不存在)
if not milvus_client.has_collection(collection_name=MILVUS_COLLECTION_NAME):
    # 创建集合
    milvus_client.create_collection(
        collection_name=MILVUS_COLLECTION_NAME,
        dimension=1536,  # text-embedding-3-small的维度是1536
        metric_type="IP",  # 内积
        index_params={"metric_type": "IP", "index_type": "IVF_FLAT", "nlist": 128}
    )
    # 插入一些测试数据
    test_products = [
        {"id": 1, "name": "Apple AirPods Pro 2", "description": "性价比高的无线蓝牙耳机,主动降噪,续航时间长", "price": 1899, "supplier": "苹果官方", "stock": 1000},
        {"id": 2, "name": "Sony WH-1000XM5", "description": "顶级主动降噪无线蓝牙耳机,音质好,续航时间长", "price": 2699, "supplier": "索尼官方", "stock": 500},
        {"id": 3, "name": "小米 Buds 4 Pro", "description": "性价比高的无线蓝牙耳机,主动降噪,音质好", "price": 899, "supplier": "小米官方", "stock": 2000},
        {"id": 4, "name": "华为 FreeBuds Pro 3", "description": "性价比高的无线蓝牙耳机,主动降噪,音质好", "price": 1299, "supplier": "华为官方", "stock": 1500},
    ]
    test_vectors = embeddings.embed_documents([p["description"] for p in test_products])
    milvus_client.insert(
        collection_name=MILVUS_COLLECTION_NAME,
        data=[{"vector": v, **p} for v, p in zip(test_vectors, test_products)]
    )
    # 创建索引
    milvus_client.create_index(
        collection_name=MILVUS_COLLECTION_NAME,
        index_params={"metric_type": "IP", "index_type": "IVF_FLAT", "nlist": 128}
    )

@tool
def vector_retrieval_tool(query: str, top_k: int = 3) -> List[Dict]:
    """
    根据用户的选品需求,从向量数据库中检索相似的历史选品数据。

    参数:
        query: 用户的选品需求
        top_k: 返回的相似历史选品数据的数量

    返回:
        相似的历史选品数据列表
    """
    try:
        # 生成查询向量
        query_vector = embeddings.embed_query(query)
        # 检索相似的历史选品数据
        results = milvus_client.search(
            collection_name=MILVUS_COLLECTION_NAME,
            data=[query_vector],
            limit=top_k,
            output_fields=["id", "name", "description", "price", "supplier", "stock"]
        )
        # 返回相似的历史选品数据
        return [r["entity"] for r in results[0]]
    except Exception as e:
        return [{"error": str(e)}]

@tool
def taobao_api_tool(product_name: str) -> Optional[Dict]:
    """
    根据产品名称,调用淘宝API获取当前的价格。

    参数:
        product_name: 产品名称

    返回:
        淘宝上的产品价格信息
    """
    try:
        # 模拟淘宝API调用
        time.sleep(random.uniform(0.1, 0.5))
        return {
            "platform": "淘宝",
            "price": random.randint(800, 2800),
            "url": f"https://s.taobao.com/search?q={product_name}",
            "status": "success"
        }
    except Exception as e:
        return {"platform": "淘宝", "error": str(e), "status": "failed"}

@tool
def jd_api_tool(product_name: str) -> Optional[Dict]:
    """
    根据产品名称,调用京东API获取当前的价格。

    参数:
        product_name: 产品名称

    返回:
        京东上的产品价格信息
    """
    try:
        # 模拟京东API调用
        time.sleep(random.uniform(0.1, 0.5))
        return {
            "platform": "京东",
            "price": random.randint(850, 2850),
            "url": f"https://search.jd.com/Search?keyword={product_name}",
            "status": "success"
        }
    except Exception as e:
        return {"platform": "京东", "error": str(e), "status": "failed"}

@tool
def pinduoduo_api_tool(product_name: str) -> Optional[Dict]:
    """
    根据产品名称,调用拼多多API获取当前的价格。

    参数:
        product_name: 产品名称

    返回:
        拼多多上的产品价格信息
    """
    try:
        # 模拟拼多多API调用
        time.sleep(random.uniform(0.1, 0.5))
        return {
            "platform": "拼多多",
            "price": random.randint(750, 2750),
            "url": f"https://mobile.yangkeduo.com/search_result.html?search_key={product_name}",
            "status": "success"
        }
    except Exception as e:
        return {"platform": "拼多多", "error": str(e), "status": "failed"}

@tool
def supplier_api_tool(product_name: str, supplier: str) -> Optional[Dict]:
    """
    根据产品名称和供应商,调用供应商API获取当前的库存。

    参数:
        product_name: 产品名称
        supplier: 供应商名称

    返回:
        供应商的库存信息
    """
    try:
        # 模拟供应商API调用
        time.sleep(random.uniform(0.1, 0.5))
        return {
            "product_name": product_name,
            "supplier": supplier,
            "stock": random.randint(0, 2000),
            "status": "success"
        }
    except Exception as e:
        return {"product_name": product_name, "supplier": supplier, "error": str(e), "status": "failed"}
3. Agent定义(src/agents.py)
from langchain_core.prompts import ChatPromptTemplate, MessagesPlaceholder
from langchain_openai import ChatOpenAI
from langchain_core.output_parsers import StrOutputParser
from typing import Dict, Any
from src.state import ProductRecommendationState
from src.tools import vector_retrieval_tool, taobao_api_tool, jd_api_tool, pinduoduo_api_tool, supplier_api_tool
from langchain.agents import AgentExecutor, create_tool_calling_agent
import os
from dotenv import load_dotenv
load_dotenv()

# 初始化LLM
llm = ChatOpenAI(model="gpt-4o-mini", temperature=0, api_key=os.getenv("OPENAI_API_KEY"))

# 选品Agent
product_selection_prompt = ChatPromptTemplate.from_messages([
    ("system", "你是一位专业的跨境电商选品专家。你的任务是根据用户的选品需求,调用向量检索工具获取相似的历史选品数据,然后生成一份详细的选品方案。选品方案应该包括产品名称、产品特点、目标用户、推荐理由等。"),
    MessagesPlaceholder(variable_name="messages"),
    ("human", "用户的选品需求:{query}\n相似的历史选品数据:{similar_products}"),
])
product_selection_tools = [vector_retrieval_tool]
product_selection_agent = create_tool_calling_agent(llm, product_selection_tools, product_selection_prompt)
product_selection_agent_executor = AgentExecutor(agent=product_selection_agent, tools=product_selection_tools, verbose=True)

def product_selection_node(state: ProductRecommendationState) -> Dict[str, Any]:
    """
    选品Agent节点

    参数:
        state: LangGraph的状态

    返回:
        更新后的状态
    """
    try:
        # 调用选品Agent
        result = product_selection_agent_executor.invoke({
            "messages": state["messages"],
            "query": state["query"],
            "similar_products": vector_retrieval_tool.invoke({"query": state["query"], "top_k": 3})
        })
        # 更新状态
        return {
            "messages": state["messages"] + [result["output"]],
            "product_selection_scheme": result["output"]
        }
    except Exception as e:
        # 更新状态
        return {
            "messages": state["messages"],
            "error": {"code": "PRODUCT_SELECTION_ERROR", "message": str(e)}
        }

# 价格对比Agent
price_comparison_prompt = ChatPromptTemplate.from_messages([
    ("system", "你是一位专业的跨境电商价格对比专家。你的任务是根据选品Agent生成的选品方案,调用淘宝API、京东API、拼多多API获取当前的价格,然后生成一份详细的价格对比报告。价格对比报告应该包括每个平台的价格、优惠信息、配送时间等。"),
    MessagesPlaceholder(variable_name="messages"),
    ("human", "选品方案:{product_selection_scheme}\n淘宝价格信息:{taobao_price}\n京东价格信息:{jd_price}\n拼多多价格信息:{pinduoduo_price}"),
])
price_comparison_tools = [taobao_api_tool, jd_api_tool, pinduoduo_api_tool]
price_comparison_agent = create_tool_calling_agent(llm, price_comparison_tools, price_comparison_prompt)
price_comparison_agent_executor = AgentExecutor(agent=price_comparison_agent, tools=price_comparison_tools, verbose=True)

def price_comparison_node(state: ProductRecommendationState) -> Dict[str, Any]:
    """
    价格对比Agent节点

    参数:
        state: LangGraph的状态

    返回:
        更新后的状态
    """
    try:
       
Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐