ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Langfuse+WebSocket构建AI对话实时监控仪表盘

Langfuse+WebSocket构建AI对话实时监控仪表盘 纯粹聊技术聊我从0到1搭这个系统时踩过的泥坑。先把话说在前面这不是一篇教程更像是我做完整个项目后的复盘笔记顺手把能“抄作业”的代码和配置都贴出来。你如果正准备给基于大模型的应用加一个实时监控仪表盘或者正被Langfuse的安装、WebSocket断连、DeepSeek回调延迟这些问题折磨这篇文章能帮你省下至少两三个晚上。整个项目用到的核心栈就五个Langfuse负责AI对话链路追踪和指标采集Langchain负责编排对话代理逻辑DeepSeek作为底层推理模型FastAPI承担后端API和WebSocket网关前端仪表盘实时渲染监控数据。下面我会沿着“为什么要这样组合”到“每一层怎么落地”的顺序展开。1. 项目背景与整体架构思路1.1 为什么需要给AI对话加监控开发完一个对话机器人之后真正难的不是让它回复而是你根本不知道它在线上到底表现如何。传统日志顶多告诉你“接口200还是500”但对话场景里用户问“帮我写个方案”模型是回答了还是胡说八道中间调了哪个工具哪些Prompt分支走了弯路这些信息普通日志完全给不出来。我之前做过一版纯日志系统把每次请求的请求体和响应体打到ELK里看起来挺全真排查问题时傻眼了日志之间没有关联ID一个对话涉及多次模型调用和工具调用时只能靠时间戳猜顺序排查一个异常案例经常要翻几百条日志。Langfuse的价值就在这里——它把一次完整的对话抽象成TraceTrace下面挂Span每次LLM调用、工具执行、外部检索都被记录成结构化事件整条调用链一目了然。另一个现实需求是实时性。对话场景的运营同学经常反馈“用户说体验很差”但没有具体案例时只能靠用户录屏。监控仪表盘直接把正在发生的对话延迟、Token消耗、异常率滚动展示出来出问题第一时间能看见不用等用户投诉。1.2 技术选型这五个组件为什么能凑到一块先看整体架构我用文字描述一下数据流用户请求进入FastAPI后端服务后端把对话内容交给Langchain构建的Agent链链路中配置DeepSeek作为LLM模型Langchain通过Langfuse的Callback Handler自动上报执行轨迹每次会话结束后后端将统计指标写入Langfuse并推送WebSocket事件仪表盘前端通过WebSocket实时接收事件并渲染图表和调用列表选Langfuse的原因很直接它开源、支持自托管、有现成的Langchain集成。Langchain生态里很多框架都内置了Langfuse的Callback比如Langchain的callbacks参数直接传langfuse_callback_handler就行不用手写埋点代码。Langchain负责的是Agent编排。DeepSeek的API接口兼容OpenAI规范直接用Langchain的ChatOpenAI类指定base_url就能接入返回的是一个标准的BaseMessage对象流和用GPT的代码路径完全一致。DeepSeek当前的API价格比主流闭源模型低不少作为监控仪表盘项目的推理引擎非常合适。FastAPI做后端没有悬念。它原生支持异步WebSocket支持得很舒服给仪表盘推送实时数据时优势很大。如果换成Flask或者DjangoWebSocket不是不能做但总要多引一个库并且处理线程模型的问题没必要。前端和WebSocket这一层是仪表盘体验的关键。最开始我想用轮询方案每5秒请求一次接口拉取最新监控数据后来在真实使用中发现两个问题一是大量无效请求打到后端浪费资源二是延迟太高对话都结束十几秒了界面才有反应。换成WebSocket之后后端事件到达即推送延迟控制在几百毫秒内。完整架构图我用一段配置描述你对照着看浏览器仪表盘 --(WebSocket/JSON)-- FastAPI服务 --(HTTP SSE回调)-- Langfuse | | (Langchain Callback) v Langchain Agent链 | v DeepSeek API2. 核心数据链路与Langfuse监控设计2.1 哪些监控指标值得收集对话监控不是什么都记记多了系统会变成垃圾场。我在设计指标时有三个原则能用一条链路串起来、能直接反映用户体验、指标必须具备可对比性。最终确定了五类核心数据会话基础信息会话ID、用户标识、开始结束时间、会话总耗时Token消耗明细每次LLM调用的输入Token、输出Token、累计Token成本延迟分布首Token延迟、完整响应延迟、工具调用延迟调用链详情每次模型调用的Prompt和Response、命中的工具名称异常信息错误类型、错误消息、失败节点位置Langfuse的数据模型天然适合这五类数据。一个Trace就是一个会话Trace内可以创建多个Span每个Span对应一次LLM调用或工具执行。我在Langchain侧做完一次问答后通过Langfuse的SDK创建Trace然后在每步工具调用、模型调用前后创建对应Span。采集方式上我做了两层。第一层是Langchain自动上报通过LangfuseCallbackHandler自动记录模型调用信息第二层是手动埋点在业务逻辑的关键位置用langfuse_client.span()记录非模型类事件比如检索器执行时间、重试次数。手动埋点虽然代码上多几行但能拿到模型调用之外的业务指标排查问题时非常有用。2.2 Langfuse自托管部署与数据模型映射如果你只是试用直接用Langfuse的云服务最快。但我们的对话数据涉及用户输入的隐私问题必须全部内网部署。Langfuse自托管还算轻松核心就是一个Docker Compose编排包含Web应用、PostgreSQL数据库和ClickHouse用于分析场景。部署完成后有两处需要立刻配置。第一处是LANGFUSE_PUBLIC_KEY和LANGFUSE_SECRET_KEYLangfuse的SDK初始化时通过这两个Key做鉴权公钥和密钥是配对的。第二处是LANGFUSE_HOST指向你部署的Langfuse服务地址。数据模型映射上我建议提前定好命名规范不然后面查询时会很混乱。我用的规范是Langfuse概念对应业务含义命名规范示例Trace完整会话conversation-{userId}-{timestamp}GenerationLLM模型调用chat-deepseek-v3Span工具/检索执行tool-search-docsEvent自定义业务动作user-feedback每个会话的Trace名称里带上用户ID和时间戳后面在Langfuse界面上按名称搜索时非常快不用翻半天。Generation的model字段固定填deepseek-chatinput和output写实际的Prompt和响应内容。3. FastAPI后端与WebSocket实时推送实现3.1 项目骨架搭建与依赖安装后端工程结构我习惯按模块划分别把路由、服务、WebSocket管理全塞到一个文件里。最终目录是这样的app/ ├── main.py ├── config.py ├── routers/ │ ├── chat.py │ └── monitor.py ├── services/ │ ├── conversation.py │ ├── langfuse_track.py │ └── ws_manager.py └── models/ ├── chat_request.py └── monitor_event.py依赖安装方面我用uv做包管理比pip快不少uv sync一下就能锁定环境。核心依赖有fastapi、uvicorn[standard]、langchain、langchain-openai、langfuse、websockets、pydantic-settings。启动命令值得注意的一点uvicorn默认在代码变更时不会自动重启开发时一定要加--reload参数。如果用了uvicorn app.main:app --reload --port 8000还是看不到热更新效果那大概率是文件路径配置问题因为uvicorn --reload默认监听当前工作目录工程代码放到了其他目录时需要配合--reload-dir参数明确指定监控目录。3.2 WebSocket连接管理器设计WebSocket服务端的核心是一个连接管理器。它要做的事情就三件维护当前所有活跃连接、向指定客户端推送消息、定时清理异常断开的连接。我写的连接管理器核心代码如下from fastapi import WebSocket from typing import List import asyncio import json class ConnectionManager: def __init__(self): self.active_connections: List[WebSocket] [] async def connect(self, websocket: WebSocket): await websocket.accept() self.active_connections.append(websocket) print(f客户端连接当前连接数{len(self.active_connections)}) def disconnect(self, websocket: WebSocket): if websocket in self.active_connections: self.active_connections.remove(websocket) print(f客户端断开当前连接数{len(self.active_connections)}) async def broadcast(self, message: dict): # 遍历副本避免在遍历过程中删除元素导致异常 for connection in self.active_connections[:]: try: await connection.send_text(json.dumps(message, ensure_asciiFalse)) except Exception: await self.disconnect(connection) manager ConnectionManager()有一点花了我不少时间排查broadcast遍历连接列表时如果某个客户端异常断开直接调用send_text会抛异常如果不处理整个推送循环会中断。我的处理方式是遍历前先复制一份连接列表推送失败就立即从活跃列表中移除该连接。同时在前端做好断线重连服务端维护好“半开连接”的兜底。WebSocket路由注册部分这样写from fastapi import APIRouter, WebSocket, WebSocketDisconnect router APIRouter() router.websocket(/ws/monitor) async def websocket_endpoint(websocket: WebSocket): await manager.connect(websocket) try: while True: data await websocket.receive_text() # 心跳回复客户端定期发送 ping服务端返回 pong if data ping: await websocket.send_text(pong) except WebSocketDisconnect: manager.disconnect(websocket)这里有个容易踩的坑receive_text()在客户端异常断开时可能不抛WebSocketDisconnect而是抛RuntimeError。我是在真实环境跑了几天之后发现有些连接断了但服务端感知不到活跃连接数越来越高最后把服务拖慢。保险的做法是加一个心跳机制客户端每30秒发一次ping服务端收到就回pong如果某条连接超过90秒没有任何消息就强制关掉并清理。3.3 Langfuse事件向WebSocket广播的联动实现后端拿到Langfuse的监控数据后怎么把数据推给前端仪表盘我最初设想了两个方案。方案一是FastAPI直接与Langfuse共用数据库前端查询时FastAPI读取Langfuse的PostgreSQL和ClickHouse数据。这个方案查询灵活但需要了解Langfuse的数据库表结构耦合度高Langfuse版本一升级表结构变了代码就得跟着改维护成本太高。方案二是Langfuse通过Webhook把事件通知到FastAPIFastAPI收到后再通过WebSocket广播给前端。后来查了下Langfuse的Webhook功能发现它提供的Webhook事件类型主要是日志导出类的粒度没有细到每次LLM调用做实时仪表盘不够灵活。最终选的是方案三也是最实用的方案不在Langfuse上做推送而是把监控数据封装到对话服务里面。在调用Langchain链的代码中我已经知道每次会话的结果、Tokens用量、延迟时间这些数据本来就会写入Langfuse同时我再构造一个轻量级事件通过WebSocket推送。Langfuse负责“事后留存和分析”WebSocket负责“实时告知”两边数据口径一致互不干扰。核心代码逻辑from models.monitor_event import MonitorEvent from services.ws_manager import manager import time async def send_monitor_event(conversation_id: str, event_type: str, payload: dict): event MonitorEvent( conversation_idconversation_id, event_typeevent_type, timestampint(time.time() * 1000), payloadpayload ) await manager.broadcast(event.model_dump())每次对话结束后把event_type设为conversation.completedpayload里放Token消耗、响应延迟、模型名称、状态码这些字段前端仪表盘收到这个事件就能实时更新。4. Langchain与DeepSeek的接入和编排实践4.1 在Langchain中配置DeepSeek模型Langchain对DeepSeek的接入方式本质上是兼容OpenAI的API协议。langchain-openai库已经封装好了流式调用、Token统计这些功能只需要把base_url指到DeepSeek的地址把模型名设为deepseek-chat就行。配置代码如下from langchain_openai import ChatOpenAI from langchain.callbacks.manager import CallbackManager from langfuse.callback import CallbackHandler langfuse_handler CallbackHandler( public_keypk-xxx, secret_keysk-xxx, hosthttps://your-langfuse-host ) llm ChatOpenAI( modeldeepseek-chat, temperature0.7, max_tokens2048, timeout60, max_retries2, base_urlhttps://api.deepseek.com/v1, api_keysk-deepseek-xxx, callbacks[langfuse_handler] )DeepSeek的base_url一定要带/v1后缀不带会被识别成非法路径。max_retries参数值得根据业务调整默认的2次重试在大并发高负载时可能触发大量重复请求如果模型本身的错误率不高可以设成1次甚至0次避免缓存了用户输入的场景里出现重复扣费。流式输出场景下还会遇到另一个问题Langchain流式返回的content是分块到达的深拷贝到callback里的数据容易只拿到最后一块。我调试时就发现Langfuse里记录的output只有一个空字符串。解决方案是在stream模式下手动把分段内容合成完整响应再交给Langfuse的generation.end()接口。如果不需要实时看到生成过程建议直接关闭流式会少很多不必要的坑。4.2 Langchain链路的Prompt与工具编排这个监控仪表盘项目里对话功能不只是简单的“一问一答”而是带工具调用和上下文记忆的Agent。我用Langchain的create_react_agent构建了一个能调用内部搜索工具的Agent。这个设计的核心思路是把“对话历史”和“工具结果”拼接成Prompt上下文。Langchain的MessagesPlaceholder可以很优雅地做到这一点不需要手工拼字符串from langchain.prompts import PromptTemplate from langchain.prompts.chat import MessagesPlaceholder prompt ChatPromptTemplate.from_messages([ (system, 你是一个智能助手可以查询内部资料库。), MessagesPlaceholder(variable_namechat_history), (human, {input}), MessagesPlaceholder(variable_nameagent_scratchpad) ])agent_scratchpad是Langchain的Agent机制必须的占位符用来存放中间推理和工具调用结果。新手容易遗漏这部分写出来的Agent在Langchain新版框架下会直接报错。工具函数写好后通过tool装饰器直接注册给Agentfrom langchain.tools import tool tool def search_docs(query: str) - str: 从内部文档库检索相关信息 # 内部检索逻辑省略 return search_result这个工具函数的docstring写得越详细越好因为Langchain会把docstring作为工具描述传给大模型让模型来决定是否调用该工具。描述含糊时模型会漏掉关键信息导致明明应该调工具的场景它偏要硬答。5. 前端仪表盘的实时展示设计5.1 仪表盘页面结构与数据聚合前端项目用的是Vue3 Vite为了减少工程负担图表库选用ECharts。仪表盘页面分成四个区域顶部统计卡片、中间对话延迟趋势图、下方Token消耗趋势图、底部实时对话记录流。顶部统计卡片展示的是“当天累计对话数”“平均响应耗时”“今日Token消耗”“异常会话数量”。这些数据由后端在WebSocket连接建立时先推送一次快照后续事件触发时增量更新前端只需要维护一个全局状态对象即可。这里分享一个从实际使用中总结的经验统计卡片上的数字更新不要过于频繁。有时对话量很大一秒内可能来十几个事件前端如果每收到一个事件就重渲染一次图表页面会很卡。我的做法是前端做一个1秒节流——事件先缓存到队列每隔1秒统一刷新一次界面。最终用户感知上几乎没有延迟但CPU占用率大幅下降。节流代码很简短let eventQueue []; let renderTimer null; function onMonitorEvent(event) { eventQueue.push(event); if (renderTimer) return; renderTimer setInterval(() { if (eventQueue.length 0) { clearInterval(renderTimer); renderTimer null; return; } processEventQueue(eventQueue.splice(0)); clearInterval(renderTimer); renderTimer null; }, 1000); }5.2 WebSocket前端连接与断线重连的处理浏览器端的WebSocket连接最让人头疼的不是初次连接而是连接中途断开后怎么办。我遇到过几种情况网络切换导致断连、服务器发布重启导致断连、连接长时间空闲被中间网络设备超时断开。前端WebSocket封装核心代码class MonitorSocket { constructor(url) { this.url url; this.ws null; this.reconnectAttempts 0; this.maxReconnectAttempts 10; this.reconnectDelay 3000; this.heartbeatTimer null; this.connect(); } connect() { this.ws new WebSocket(this.url); this.ws.onopen () { console.log(监控连接已建立); this.reconnectAttempts 0; this.startHeartbeat(); }; this.ws.onmessage (event) { // 心跳响应处理 if (event.data pong) return; const data JSON.parse(event.data); this.handleMonitorEvent(data); }; this.ws.onclose () { console.log(监控连接已关闭); this.stopHeartbeat(); this.reconnect(); }; this.ws.onerror (error) { console.error(WebSocket错误, error); this.ws.close(); }; } reconnect() { if (this.reconnectAttempts this.maxReconnectAttempts) { console.error(重连次数过多停止重连); return; } const delay Math.min(this.reconnectDelay * Math.pow(2, this.reconnectAttempts), 30000); console.log(等待 ${delay}ms 后重连...); setTimeout(() { this.reconnectAttempts; this.connect(); }, delay); } startHeartbeat() { this.heartbeatTimer setInterval(() { if (this.ws.readyState WebSocket.OPEN) { this.ws.send(ping); } }, 30000); } stopHeartbeat() { if (this.heartbeatTimer) { clearInterval(this.heartbeatTimer); this.heartbeatTimer null; } } }断线重连有两个细节必须说。第一是重连延迟要使用指数退避策略初始3秒失败一次翻一倍最大不超过30秒。否则服务器恢复后所有客户端的重连请求同时到达直接把服务打挂。第二是重连成功后前端要主动拉取一次全量快照恢复断线期间丢失的事件避免图表出现数据缺口。我在实际开发中还遇到过前端H5页面能正常连接WebSocket但打包成App后连不上的情况。排查了半天发现是因为App里访问的远端地址和H5测试地址不一致安全域名限制导致连接被拒。如果你的移动端也遇到了类似问题优先检查网络权限和域名白名单。6. 部署、常见问题与实战心得6.1 Langfuse部署与接入的坑Langfuse自托管部署成功后最容易出现的问题是SDK连接不上服务端。排查步骤我整理成了表格现象排查方向解决方案SDK报401认证失败公钥/私钥配置错误或粘贴了多余空格检查LANGFUSE_PUBLIC_KEY和LANGFUSE_SECRET_KEY是否配对SDK报404路径不存在base_url或host路径配置错误确认host末尾不带/路径保持/api/public仪表盘没有数据回调Handler未挂载到Langchain链检查callbacks[langfuse_handler]是否传入了链和LLM数据延迟到达SDK采用异步批量上报调整LANGFUSE_PROCESSOR_BATCH_SIZE和刷新间隔我在第一次接入时犯过一个低级错误把Langfuse的public_key和secret_key顺序写反了调试了半小时一直在怀疑网络问题。大家配置时注意一下这两个Key在官方文档的位置别想当然。6.2 WebSocket断连、性能优化与稳定性经验WebSocket的1006错误码是我被问得最多的问题。这个错误码表示连接非正常关闭但WebSocket协议没有给出具体原因。排查方向通常有三个一是服务器重启连接被动断开二是前端页面有多个连接部分连接被聊天窗口复用导致冲突三是服务器Nginx或网关配置了空闲超时长时间没有消息交互就会断。针对最后一个原因解决方案是在Nginx层调整WebSocket的超时参数proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; proxy_read_timeout 3600s; proxy_send_timeout 3600s;proxy_read_timeout表示读取上游数据的超时时间默认60秒如果WebSocket在60秒内没有数据交互连接就会被Nginx切断。调整到3600秒后配合每30秒一次的心跳连接稳定性明显提升。性能优化方面WebSocket广播在连接数量大时会有性能隐患。broadcast循环是串行发送如果某个客户端网络慢会阻塞后面所有客户端的发送。优化方案有两个方向一是把发送操作改为异步任务交给asyncio.create_task并发执行二是使用发布订阅模式把WebSocket连接挂到Redis频道上横向扩容时也能正常工作。当前项目连接数在几十到几百之间串行发送完全够用但如果未来连接数上千就得考虑改造了。还有一个被忽略的细节是消息体的大小。有段时间我发现仪表盘页面越来越卡打开开发者工具一看WebSocket推送的单个消息居然有好几MB因为在payload里把完整的Prompt和Response都塞进去了。仪表盘的实时视图不需要完整文本只要摘要和统计信息。我把推送内容改成了截断后的摘要完整数据仍然通过Langfuse存储前端只展示精简事件页面一下流畅了很多。6.3 项目最终效果与扩展建议整套系统上线后效果非常直观。运营侧的同事现在可以实时看到每一条对话的状态和Token消耗后端同学通过Langfuse的界面可以下钻到任意一条Trace查看完整的模型输入输出。有一次线上对话突然大面积超时仪表盘上的“平均响应耗时”曲线提前两分钟就开始走高我们比用户投诉更早发现了问题定位到是DeepSeek侧出现了延迟抖动立刻做了降级切换。后续扩展我有两个方向规划。第一个方向是把监控指标接入告警当某个时间窗口内Token消耗异常增长或模型错误率超过阈值时自动通过Webhook通知到企业微信群。第二个方向是把用户反馈引入监控体系在对话结束时让用户对回复质量打分评分数据与对话Trace关联找出高频低分场景。如果你打算在生产环境落地这个架构我还想提醒一点监控本身也有成本。Langfuse会保存大量完整对话记录数据量增长很快需要定期清理归档。建议在初始化SDK时配置采样率线上高并发场景采样率设为10%到50%就够了不需要所有请求都进监控系统否则存储成本和查询性能都会成为新问题。做这套系统的最大体会是AI应用的可观测性远没有现成工具可以一步到位但把Langfuse的追踪能力、Langchain的回调机制和WebSocket的实时通道组合起来确实能够搭出接近商业产品的效果。希望这些实践经验能帮你少走一些弯路。
返回列表