ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

AI对话监控仪表盘全链路搭建:Langfuse+Langchain+FastAPI实践

AI对话监控仪表盘全链路搭建:Langfuse+Langchain+FastAPI实践 这几年做大模型应用我踩过最深的坑不是模型效果不好而是不知道它在生产环境里到底发生了什么。用户报一个“回答变慢了”“突然报错了”你连是模型接口超时还是 Prompt 构造出了问题都分不清。后来我搭了一套 AI 对话监控仪表盘把 Langfuse、Langchain、DeepSeek、FastAPI 和 WebSocket 串成一条完整的实时链路才算真正对线上对话有了掌控感。这篇文章把这个项目的完整思路和源码细节拆开讲包括为什么这么选型、每一层怎么对接、以及我在实操中遇到的坑希望能给你一个可以直接对照复现的参考。这套方案适合谁如果你正在做基于大模型的对话产品或者准备上线一个带流式输出的 AI 应用想知道每次调用花了多少钱、消耗了多少 Token、延迟怎么样、用户到底满不满意那这套监控方案就是为你准备的。Langfuse 负责链路追踪和指标记录Langchain 负责编排模型调用DeepSeek 作为底层模型FastAPI 提供接口服务WebSocket 把对话和监控数据实时推到前端仪表盘整个链路清晰、可扩展。1. 整体设计思路为什么是这五个组件1.1 先搞清楚要监控什么在选型之前我先把需求列了个清单。所谓 AI 对话监控不只是把聊天记录存下来这么简单至少要有这几个维度成本监控每次对话消耗多少 TokenDeepSeek 的输入输出 Token 分别多少累计花费多少。性能监控请求延迟首 Token 延迟总耗时。质量监控用户的问题是什么模型的回答是什么是否需要人工介入。异常监控哪些请求失败了失败原因是什么是超时、限流还是内容审核拦截。实时性对话过程中的中间状态能实时看到而不是等整个请求结束才回看日志。这些需求决定了架构必须是“追踪 业务 展示”三层解耦的。追踪层记录模型调用的全链路轨迹业务层负责对话逻辑和服务暴露展示层负责把数据实时渲染到界面。1.2 技术选型的核心考量选 Langfuse 当时也犹豫过是不是太重了毕竟市面上还有 Phoenix、MLflow 之类的方案。但我最终选它有几个很实在的理由。Langfuse 专门为 LLM 应用设计能原生记录 Prompt、Token 数、模型参数、延迟这些指标而且提供了 LLM 行业的“标准格式” Trace 和 Span 概念。配合 Langchain 有现成的 Callback Handler接入成本极低。Langfuse 可以自托管数据不出内网这对很多公司的合规要求来说很重要。它的看板自带成本分析和评分功能省去自己从零做统计报表的功夫。Langchain 我承认它有点重早期版本 API 变动也频繁但对于这种多步骤对话场景它的好处是提供了统一的模型抽象换模型不用改业务代码有现成的回调机制便于埋点生态里有很多组件比如记忆、检索都可以复用。DeepSeek 作为模型层看中的是它的性价比和国产模型的稳定性。API 兼容 OpenAI 格式接入 Langchain 非常顺滑。在监控这件事上它还能额外记录一些模型特有的返回字段。FastAPI 不用多说异步原生支持WebSocket 支持得很干净做实时转发非常合适。Python 生态里跟 Langchain 配合最舒服的也是它。1.3 整体架构与数据流向整个系统的数据流是这样的用户在前端页面发起对话消息通过 WebSocket 发送到 FastAPI 后端。FastAPI 拿到消息后调用 Langchain 编排的链或 AgentLangchain 内部调用 DeepSeek 模型。在调用过程中Langfuse 的 Callback Handler 会自动把这一轮调用的 Prompt、Token、延迟、成本等数据写入 Langfuse 的数据库。与此同时FastAPI 把模型流式返回的 Token 通过同一个 WebSocket 推回前端展示监控侧的数据则通过另一个 WebSocket 通道推送到仪表盘页面实时刷新指标卡片和 Trajectory 列表。这里我特意把对话消息和监控消息分成两个 WebSocket 通道而不是混在一个通道里。原因是对话流对实时性要求高如果被监控数据的批量推送挤占带宽用户会感觉到打字机效果一顿一顿的。分开走互不干扰排查问题也方便。2. Langfuse 部署与 Langchain 接入细节2.1 自托管部署 LangfuseLangfuse 官方提供 Docker Compose 方式我建议直接用它不用手动去装 PostgreSQL 和 ClickHouse它会把依赖一起编排好。# docker-compose.yml version: 3.9 services: langfuse: image: langfuse/langfuse:2 restart: unless-stopped depends_on: db: condition: service_healthy environment: DATABASE_URL: postgresql://langfuse:langfusedb:5432/langfuse NEXTAUTH_URL: http://localhost:3000 NEXTAUTH_SECRET: mysecretkey SALT: mysalt ENCRYPTION_KEY: myencryptionkey REDIS_URL: redis://redis:6379 ports: - 3000:3000 db: image: postgres:16 restart: unless-stopped environment: POSTGRES_USER: langfuse POSTGRES_PASSWORD: langfuse POSTGRES_DB: langfuse volumes: - langfuse_db:/var/lib/postgresql/data healthcheck: test: [CMD-SHELL, pg_isready -U langfuse] interval: 5s timeout: 5s retries: 5 redis: image: redis:7 restart: unless-stopped volumes: langfuse_db:启动之后访问http://localhost:3000注册管理员账号然后创建一个项目拿到 API Key公钥和私钥。这两个 Key 后面要配置到 Langchain 的回调里。注意Langfuse v2 和 v3 的环境变量略有不同。如果你用的是 V4 以后的版本Redis 和 ClickHouse 几乎是必须的配置项也会多出CLICKHOUSE_URL、CLICKHOUSE_DATABASE这些。我这里给的是 v2/v3 兼容的写法实际以你拉取的镜像版本为准。2.2 在 Langchain 中接入 DeepSeekDeepSeek 的 API 兼容 OpenAI 格式所以在 Langchain 里不需要特殊的包装器直接使用ChatOpenAI并修改 base_url 和 api_key 就行。from langchain_openai import ChatOpenAI llm ChatOpenAI( modeldeepseek-chat, temperature0.7, api_keysk-xxxxxxxx, base_urlhttps://api.deepseek.com/v1, streamingTrue, )这里有个坑就是streamingTrue必须显式打开。如果你不打开模型会等完整回复生成完毕才返回WebSocket 的“打字机”效果就实现不了。同时流式模式下Langchain 的stream方法会逐块产出 Token而 Langfuse 的生成追踪需要依赖这些 Token 累积来计算 Token 数。2.3 Langfuse 回调接入一行代码开启全链路追踪Langchain 接入 Langfuse 做监控官方提供的是回调机制。在 Langchain 生态里只需要引入CallbackHandler然后传给链或 LLM 调用即可。from langfuse.callback import CallbackHandler langfuse_handler CallbackHandler( public_keypk-xxx, secret_keysk-xxx, hosthttp://localhost:3000 ) # 在调用链时传入回调 result chain.invoke( {question: user_message}, config{callbacks: [langfuse_handler]} )用config{callbacks: [handler]}这种方式而不是在构建链的时候传入原因是为了每个请求都能正确生成独立的 Trace。如果你把 CallbackHandler 绑定在链对象上多个用户并发时 trace_id 会串监控数据就乱了。如果你用langchain_core.runnables.RunnablePassthrough做数据流传递同样可以传 callbacks。记住一个原则Langfuse 的 Trace 是跟着一次完整调用走的只要回调在调用入口处传入里面任何一步的 Model 调用都会被自动追踪为 Span。2.4 LangGraph 和 Langchain 的关系项目初期我也纠结要不要直接上 LangGraph。Langchain 是基础框架LangGraph 是在它之上的图编排引擎适合有状态、有分支的 Agent 流程。如果只是简单问答链路用 Langchain 就够了如果以后要做多工具编排、条件跳转、人工审批循环那时候再把核心节点迁移到 LangGraph 也不晚。这不是二选一的问题LangGraph 底层还是要用 Langchain 的模型封装和工具抽象。3. FastAPI 后端与 WebSocket 实时通信构建3.1 项目结构划分后端我按照“路由、服务、流控”三层拆分不复杂但边界清楚。app/ ├── main.py # FastAPI 入口 ├── routers/ │ ├── chat.py # 对话相关路由含 WebSocket │ └── dashboard.py # 监控数据推送路由 ├── services/ │ ├── llm_service.py # 封装 Langchain 调用 │ └── trace_service.py # 查询 Langfuse 数据 └── ws_manager.py # WebSocket 连接管理器3.2 实现一个连接管理器WebSocket 管理最容易出问题的是并发连接多了之后不知道该往哪个连接发数据。我写了一个简单的管理器维护连接池并提供广播和单发能力。# ws_manager.py import asyncio from typing import Dict, Set from fastapi import WebSocket class ConnectionManager: def __init__(self): self.active_connections: Dict[str, Set[WebSocket]] {} async def connect(self, room: str, websocket: WebSocket): await websocket.accept() if room not in self.active_connections: self.active_connections[room] set() self.active_connections[room].add(websocket) def disconnect(self, room: str, websocket: WebSocket): if room in self.active_connections: self.active_connections[room].discard(websocket) if not self.active_connections[room]: del self.active_connections[room] async def send_to_room(self, room: str, message: dict): if room not in self.active_connections: return for conn in list(self.active_connections[room]): try: await conn.send_json(message) except Exception: await self.disconnect(room, conn)这里用set存储连接避免重复连接时误加常量。发送时做异常捕获防止个别连接断了之后整条广播链路被阻塞。3.3 对话 WebSocket 端点对话端点负责接收用户消息、调用 Langchain、把流式结果推回前端。核心逻辑如下# routers/chat.py from fastapi import APIRouter, WebSocket, WebSocketDisconnect from services.llm_service import stream_chat from ws_manager import ConnectionManager import json router APIRouter() manager ConnectionManager() router.websocket(/ws/chat) async def chat_endpoint(websocket: WebSocket): await manager.connect(chat, websocket) try: while True: data await websocket.receive_text() payload json.loads(data) user_message payload.get(message, ) # 立即回执让前端显示用户消息已发出 await websocket.send_json({type: status, data: received}) # 流式调用 Langchain async for chunk in stream_chat(user_message): await websocket.send_json({ type: token, data: chunk.content }) # 发送这轮对话的结束标记 await websocket.send_json({type: end, data: }) except WebSocketDisconnect: manager.disconnect(chat, websocket)这里有个容易被忽略的点while True接收消息实现多轮对话但每个用户的多次提问会共享同一个 WebSocket。如果有人反复快速提问后一个请求还没返回前一个流式输出还在推前端就乱套了。我加了一个简单的并发控制在前端用户发送消息后禁用发送按钮收到end之后才恢复。这在单用户场景下足够如果要支持多并发对话就得上消息 ID 配对复杂不少。3.4 模型流式调用服务stream_chat这层是 Langchain 和 DeepSeek 衔接的地方。# services/llm_service.py from langchain_openai import ChatOpenAI from langchain.schema import HumanMessage from langfuse.callback import CallbackHandler llm ChatOpenAI( modeldeepseek-chat, base_urlhttps://api.deepseek.com/v1, api_keysk-xxx, streamingTrue, ) langfuse_handler CallbackHandler( public_keypk-xxx, secret_keysk-xxx, hosthttp://localhost:3000 ) async def stream_chat(user_message: str): async for chunk in llm.astream( [HumanMessage(contentuser_message)], config{callbacks: [langfuse_handler]} ): yield chunk这里直接用llm.astream因为 FastAPI 是异步的如果还用同步的stream会在模型等待响应时阻塞整个事件循环。这是我的一个经验教训一开始我用同步的chain.stream包在run_in_executor里非常别扭后来发现 Langchain 本身已经提供了astream直接用它就行。另外注意一点Langfuse 的回调和astream配合时Trace 是在生成结束后才完整落库的。流式过程中前端能看到 Token但 Langfuse 看板上的指标要等这一轮生成完成才刷出来。这符合常理毕竟 Token 数和成本只有生成全部结束才知道。4. 仪表盘数据通道把 Langfuse 指标实时推到前端4.1 监控数据从哪来仪表盘要显示的内容分两类。一类是 Langfuse 已经统计好的指标比如总请求次数、平均延迟、累计成本、Token 消耗这类数据可以通过 Langfuse 的 API 定时拉取。另一类是实时的对话流这要求后端在每次模型调用时主动推送。我采用的是混合方案对话内容流走 WebSocket 实时推指标卡片每 5 秒从 Langfuse API 拉取一次。长轮询还是 WebSocket指标卡片的刷新频率要求不高定时拉取更简单稳定对话内容则必须用 WebSocket因为中间有“打字机”式输出轮询无法还原逐 Token 的效果。4.2 拉取 Langfuse 指标Langfuse 有成熟的 Python SDK可以直接查询 Trace 和观察数据。# services/trace_service.py from langfuse import Langfuse langfuse Langfuse( public_keypk-xxx, secret_keysk-xxx, hosthttp://localhost:3000 ) def get_recent_metrics(minutes: int 10): traces langfuse.fetch_traces(limit50) total_tokens 0 total_cost 0.0 total_latency 0.0 for trace in traces.data: usage trace.usage or {} total_tokens usage.get(totalTokens, 0) total_cost trace.cost or 0 total_latency trace.latency or 0 count len(traces.data) return { count: count, total_tokens: total_tokens, total_cost: round(total_cost, 4), avg_latency: round(total_latency / count, 2) if count else 0 }这里trace.cost是 Langfuse 根据模型和 token 估算的DeepSeek 的价格结构在 Langfuse 里如果没内置可能会被算成 0 或者不准。我是手动在 Langfuse 后台的模型配置里添加了deepseek-chat的计价规则这样成本数据才有参考价值。4.3 仪表盘 WebSocket 推送仪表盘页面的连接地址是/ws/dashboard。后端每隔 5 秒拉取一次指标并推送给所有已连接的仪表盘页面。# routers/dashboard.py import asyncio from fastapi import APIRouter, WebSocket, WebSocketDisconnect from ws_manager import ConnectionManager from services.trace_service import get_recent_metrics router APIRouter() manager ConnectionManager() router.websocket(/ws/dashboard) async def dashboard_endpoint(websocket: WebSocket): await manager.connect(dashboard, websocket) try: while True: metrics get_recent_metrics() await manager.send_to_room(dashboard, { type: metrics, data: metrics }) await asyncio.sleep(5) except WebSocketDisconnect: manager.disconnect(dashboard, websocket)这个循环里每次send_to_room之后必须sleep。如果不 sleep前端连接断开时这个循环会疯狂重连空转把 CPU 打满。我还把get_recent_metrics的异常包在 try 里了Langfuse 接口偶尔闪断不能让仪表盘整个挂掉。4.4 前端仪表盘页面仪表盘我用纯 HTML 原生 JS 写的没有上 React因为监控页面的交互密度不高原生实现更快也方便你嵌入到现有的运维系统里。!DOCTYPE html html head meta charsetutf-8 titleAI 对话监控/title style body { font-family: -apple-system, PingFang SC, sans-serif; padding: 24px; } .cards { display: flex; gap: 16px; margin-bottom: 20px; } .card { border: 1px solid #e5e7eb; border-radius: 8px; padding: 16px; min-width: 180px; } .card .value { font-size: 32px; font-weight: 600; } #log { background: #0f172a; color: #e2e8f0; padding: 16px; border-radius: 8px; height: 400px; overflow-y: auto; font-family: monospace; } /style /head body h3实时监控面板/h3 div classcards div classcard请求数span idcount classvalue0/span/div div classcardToken 总量span idtokens classvalue0/span/div div classcard成本(USD)span idcost classvalue0/span/div div classcard平均延迟(秒)span idlatency classvalue0/span/div /div div idlog/div script const ws new WebSocket(ws://localhost:8000/ws/dashboard); ws.onmessage (event) { const msg JSON.parse(event.data); if (msg.type metrics) { document.getElementById(count).textContent msg.data.count; document.getElementById(tokens).textContent msg.data.total_tokens; document.getElementById(cost).textContent msg.data.total_cost; document.getElementById(latency).textContent msg.data.avg_latency; } }; ws.onclose () console.warn(dashboard ws closed); /script /body /html前端要重点处理断线重连。浏览器在 Wi-Fi 切换或者服务端重启时 WebSocket 会迅速断开如果没做重连监控页面就变成一张死图。我一般用简单的重试策略function connect() { const ws new WebSocket(ws://localhost:8000/ws/dashboard); ws.onclose () setTimeout(connect, 3000); } connect();重连间隔 3 秒比较合适太短容易在服务端重启过程中疯狂握手把日志刷爆太长又会让监控空窗太久。5. 对话页面与流式展示的完整闭环5.1 对话页面的 WebSocket 交互对话页面相对简单但它和仪表盘共用一套后端只是 room 不同。我直接写了一个支持多轮消息的聊天窗口。const chatSocket new WebSocket(ws://localhost:8000/ws/chat); function sendMessage() { const input document.getElementById(input); const message input.value.trim(); if (!message) return; appendMessage(user, message); chatSocket.send(JSON.stringify({ message })); input.value ; document.getElementById(sendBtn).disabled true; } chatSocket.onmessage (event) { const msg JSON.parse(event.data); if (msg.type token) { appendToken(msg.data); } else if (msg.type end) { document.getElementById(sendBtn).disabled false; } };appendToken的实现方式有讲究。如果你每次都对整个回答区域做innerHTML 老内容 新 token在 Token 频率高的时候浏览器会卡顿。我采用的方式是维护一个游标累积到一个微批次再更新。或者直接创建一个span引用用textContent 的方式追加。let answerSpan null; function appendToken(token) { if (!answerSpan) { answerSpan document.createElement(span); document.getElementById(answer).appendChild(answerSpan); } answerSpan.textContent token; }5.2 对话流和监控流的联动对运维来说看着对话页面和监控页面同时刷新才能体会到这套系统的价值。用户每发出一个问题监控面板的请求数加一Token 总量跳动等回答打完字延迟数据出现成本数字更新。要调试这种联动可以在 Langfuse 的 Trace 详情里点开这个请求查看完整的 Prompt 和每个步骤的耗时。这样线上问题的定位从“猜”变成了“看”。我还做了一个小功能在仪表盘的日志区展示最新几条对话摘要。实现方式是在stream_chat里按用户消息缓存一份简短的记录FastAPI 的监控推送循环把它一并带过去。这个可以自由扩展比如加一个“违规内容触发次数”的计数器。5.3 整合入口最后在main.py里挂上两个路由# main.py from fastapi import FastAPI from fastapi.staticfiles import StaticFiles from routers import chat, dashboard app FastAPI(titleAI 对话监控仪表盘) app.include_router(chat.router) app.include_router(dashboard.router) app.mount(/, StaticFiles(directorystatic, htmlTrue), namestatic)注意app.mount(/, StaticFiles(...))必须放在最后否则它会拦截所有路由导致/ws/chat和/ws/dashboard根本进不去。这是我第一次跑起来发现 WebSocket 一直 404 的原因坑很深。6. 常见问题与排查技巧实录6.1 已遇到的问题速查表我在搭这个项目的过程中遇到的问题比预想多整理成一张表你能直接对照排查。现象可能原因解决办法Langfuse 看板没有数据CallbackHandler 没有传在调用入口检查config{callbacks: [handler]}是否在 invoke 时传入WebSocket 连接反复断开 code 1006服务端异常导致连接被清理用浏览器 Network 面板看握手和错误信息检查后端日志DeepSeek 响应慢但界面无感知没有开启 streaming设置streamingTrue并确认使用astream指标卡片成本为 0Langfuse 没有内置 DeepSeek 计价在 Langfuse 模型配置里手动添加 deepseek-chat 价格对话页面 token 显示重复前端对同一条数据流接收了多次检查是否在 onmessage 里多次 append或重连后未清空旧消息/路径路由被静态文件劫持mount 顺序不对把app.mount(/, StaticFiles(...))放到最后仪表盘光标卡顿每次 token 都操作 DOM累积 token 批量更新或使用 textContent 追加设置了回调但 trace 里看不到 token 数模型未开启流式但 handler 未拿到 usage升级 langchain-openai 版本检查 DeepSeek 返回的 usage 字段6.2 调试 WebSocket 的三种手段WebSocket 调试比普通 HTTP 麻烦我给你三个实用的手段。第一浏览器开发者工具的 Network 标签页里选择 WS能看到实时帧排查握手和数据收发最方便。第二用 Python 的websockets库写一个最小客户端排除浏览器干扰。第三后端日志里打印每次连接和断开的事件配一个on_disconnect回调观察连接生命周期。我在排查 code 1006 时就是靠后端日志发现是 Langfuse 回调在主线程里阻塞了事件循环导致 WebSocket 心跳超时被底层判定为断开。解决办法是把 Langfuse 回调的落库动作改成异步方式或者至少不要放在高并发的 WebSocket 收发路径上做同步写操作。6.3 关于 Langfuse 版本升级的细节Langfuse 从 v2 升到 v3 甚至 v4最大的变化是数据存储从纯 Postgres 迁移到 ClickHouse查询性能提升明显但部署复杂度也上去了。如果你只是个人项目或者小团队内部用v3 的 Docker Compose 部署完全够如果要支撑大规模生产查询建议直接上 v4并配置好 ClickHouse 的持久化。升级前记得备份 Postgres 数据尤其是 trace 和 observation 表。Langfuse 官方提供迁移脚本但我在实际操作中发现迁移后部分旧 trace 的 cost 字段会丢失需要重新拉取一次历史数据才能补全。7. 额外的心得让监控体系真正发挥价值搭建完这套系统我有几点体会比较深。不要迷信仪表盘上的数字先确认数据链路是完整的。很多团队把监控搭起来后发现看板很漂亮但关键请求根本没被记录这种监控反而误导人。我的习惯是先用一个固定 ID 跑一条测试对话然后在 Langfuse 里直接按 trace_id 搜到它确认链路通了再去看聚合指标。Langfuse 的 Trace 结构是和业务绑定的。你可以给每条 trace 添加自定义的session_id、user_id、input和output元数据这样后续做用户维度的分析会方便很多。我在调用链里额外传了用户的会话 ID 和客户端 IP 的脱敏哈希这样能在 Langfuse 里筛选某个用户的全部对话。监控这件事宁可先粗后细不能一开始就求全。第一版先把请求数、Token、延迟、报错率这四个核心指标做出来再逐步加上质量评分、人工标注、Prompt 版本对比这些高级功能。否则项目很容易在数据模型设计里陷进去迟迟上不了线。最后说一个小技巧在 FastAPI 里给 WebSocket 加一个简单的鉴权比如让前端在连接时带上一个 token query 参数后端验证通过才 accept。虽然 WebSocket 的鉴权不能替代 HTTP 的完整登录体系但至少能挡住误连和恶意扫描。别裸奔线上环境的扫描器比你想象的多这是我在把服务暴露到公网后收到的教训。
返回列表