ARTICLE DETAIL

资讯详情

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

基于Langfuse和WebSocket的AI对话监控仪表盘搭建实战

基于Langfuse和WebSocket的AI对话监控仪表盘搭建实战 从去年开始我团队里接入的大模型应用越来越多从客服问答到内部知识库再到各种RAG Agent。代码写起来很快但一上线就头疼用户发一句话过来链路里到底走了哪个模型、传了什么上下文、花了多少token、哪个环节响应变慢了全是一团黑。出了问题只能让人复现再一层层打日志效率低得让人想摔键盘。后来我把Langfuse、Langchain、DeepSeek、FastAPI和WebSocket这套组合搭成了一个实时的AI对话监控仪表盘才算把这口黑锅掀开了一角。这篇文章就完整记录一下我是怎么从零把它搭起来的包括每一步的关键代码、选型原因和实际踩过的坑希望能给同样在做大模型应用观测和实时链路追踪的朋友一点参考。1. 项目背景与整体架构思路1.1 为什么需要一套自建的AI对话监控系统先说说我为什么要做这个事。市面上的APM工具不少但大多数是针对传统Web服务的能看懂HTTP状态码、数据库慢查询、JVM指标。到了LLM应用这个层面监控的对象完全变了你关心的是每次对话里提示词怎么组装、模型的temperature和top_p参数是多少、输入了多少字符、输出了多少token、首字延迟多久、整轮响应耗时多少以及最关键的一环——中间有没有触发工具调用或者检索逻辑。这些信息传统日志系统很难结构化地记录下来。我最初的做法是手动在代码里埋点把每次请求的输入输出打到日志文件然后去Kibana里查。这样能用但体验很差。一个是请求多了以后日志文件膨胀得厉害检索一次要等很久另一个是没法做链路追踪一个对话里可能调了多次模型它们之间的父子关系完全靠猜。后来我研究了一圈决定把Langfuse作为观测核心因为它是专门为LLM应用设计的开源可观测平台可以记录trace、span、generation、token消耗和评分功能正好打在痛点上。1.2 技术栈选型与各组件角色既然要自建技术栈就不能乱选。我最终确定的组合是Langfuse负责数据采集与可视化Langchain负责编排大模型调用链路DeepSeek作为实际推理的大模型底座FastAPI提供后端API服务和WebSocket长连接前端再配一个轻量的HTML页面实时展示监控数据。这套组合里每个组件都有明确的分工不是随便拼凑的。Langfuse开源LLM可观测平台支持自托管用Docker Compose就能跑起来。它提供Python SDK和Langchain Callback集成可以在不入侵业务代码的前提下自动捕获模型的输入输出、耗时、token消耗并生成可视化的链路图。Langchain负责把大模型调用封装成标准接口并管理提示词模板、工具调用、记忆模块等。它和Langfuse有官方集成能自动把链路数据同步过去。DeepSeek底层大模型服务提供OpenAI兼容的API接口可以无缝接入Langchain的ChatOpenAI封装省去自己写HTTP调用的麻烦。FastAPI基于Python的异步Web框架天然支持WebSocket和异步任务适合做实时数据中转层。它同时承担两个职责把前端的对话请求转发给Langchain链路同时把过程中产生的token文本、耗时指标实时推送到监控页面。WebSocket实现服务端到浏览器的全双工通信解决监控数据实时刷新的问题。这套架构里最关键的一条原则是业务对话通道和监控通道要解耦。用户的对话请求走正常的API调用流程通过Langchain处理监控数据则由Langfuse在后端采集再通过另一个WebSocket通道推给前端展示。二者互不阻塞即使监控系统挂了也不影响业务主流程。1.3 数据链路与实时性设计整体数据流向是这样的用户在监控页面输入一条消息页面通过WebSocket把消息发给FastAPI后端后端用Langchain组装好提示词调用DeepSeek的API模型流式返回token的过程中Langchain通过其Callback机制把每个token、当前阶段的耗时、token累计消耗量等事件实时传给Langfuse的Callback Handler同时我们自己写一个自定义Callback把这些事件转发给WebSocket ManagerWebSocket Manager把数据推送到浏览器端仪表盘实时更新token文本、延迟曲线和成本统计。这里有一个细节值得说明为什么不直接让前端轮询Langfuse的API来刷新数据因为Langfuse的API是有延迟的trace数据入库后通常要等几秒才能查到而且轮询的效率很低还会增加Langfuse服务端的压力。用WebSocket做实时推送事件到了就推到了就渲染延迟基本在100ms以内体验完全不一样。2. 环境准备与前期部署要点2.1 DeepSeek接入OpenAI兼容接口的妙处DeepSeek是我这次选的模型服务。它提供了OpenAI兼容的API所以接入Langchain时不需要任何特殊适配直接用ChatOpenAI这个类把base_url改成DeepSeek的地址就行。这个兼容性设计真的节省了大量时间如果每个模型服务都要写一套独立的SDK封装项目复杂度会翻好几倍。from langchain_openai import ChatOpenAI llm ChatOpenAI( modeldeepseek-chat, api_keysk-你的DeepSeek密钥, base_urlhttps://api.deepseek.com, temperature0.7, streamingTrue, )需要注意的是DeepSeek的上下文长度和计费方式和OpenAI不完全一样在代码里要明确把streaming参数打开。因为监控仪表盘要展示流式输出效果如果用非流式接口前端要等模型全部输出完才能拿到结果实时性就没了。另外DeepSeek的API Key管理后台可以创建多个Key建议为这个监控项目单独申请一个方便排查问题和控制预算。2.2 Langfuse自托管Docker Compose一键部署Langfuse官方提供了自托管方案。我用Docker Compose在服务器上部署了一套整个流程比较顺畅。在服务器上创建一个目录写一个docker-compose.yml文件内容大致是启动Langfuse的Web服务、PostgreSQL数据库和Redis缓存然后执行docker compose up -d等镜像拉取完成访问3000端口就能看到登录页面。首次启动后需要注册管理员账号然后创建一个Project拿到公钥和私钥这两个Key后面集成时要用。version: 3.9 services: langfuse: image: langfuse/langfuse:2 ports: - 3000:3000 environment: DATABASE_URL: postgresql://postgres:postgresdb:5432/langfuse NEXTAUTH_URL: http://localhost:3000 NEXTAUTH_SECRET: mysecret SALT: mysalt ENCRYPTION_KEY: myencryptionkey depends_on: - db db: image: postgres:15 environment: POSTGRES_USER: postgres POSTGRES_PASSWORD: postgres POSTGRES_DB: langfuse volumes: - postgres_data:/var/lib/postgresql/data redis: image: redis:7 volumes: postgres_data:版本方面目前的Langfuse v3/v4系列在控制台界面和API上有些差异但基础的SDK用法和Callback Handler保持兼容。我这套是基于v2镜像部署的如果你拉的是最新版本可能在登录页面和项目设置的位置上略有不同但整体逻辑一致。2.3 Langchain项目结构与依赖版本坑Langchain这个库的版本迭代非常快不同版本之间的API差异很大。早期的Langchain把所有东西都塞进langchain这一个包里后来拆成了langchain-core、langchain-community、langchain-openai等多个子包。我在项目里用的版本是Langchain 0.1.x系列配合langchain-openai 0.1.x整体比较稳定。项目目录结构我建议这样组织ai-monitor-dashboard/ ├── main.py # FastAPI入口 ├── config.py # 全局配置项 ├── callbacks.py # 自定义Callback实现 ├── langchain_app.py # Langchain链路构建 ├── ws_manager.py # WebSocket连接管理器 ├── requirements.txt └── templates/ └── index.html # 前端监控页面依赖清单如下fastapi uvicorn[standard] langchain0.1.* langchain-openai0.1.* langfuse2.* websockets jinja2 python-dotenv这里特别提醒一点不要无脑装最新版。Langchain的API变动非常大比如runnable的invoke方法、CallbackHandler的接口都可能在新版本里改名。固定版本号是防止项目第二天就起不来的关键措施。2.4 Langfuse环境变量集成Langchain和Langfuse的集成有两种方式。第一种是通过环境变量自动集成只要设置好LANGCHAIN_TRACING_V2和Langfuse相关环境变量Langchain的调用链就会自动上报到Langfuse。第二种是手动创建Langfuse CallbackHandler然后在调用chain时显式传入callbacks配置。我推荐第二种方式。原因是环境变量集成虽然省事但可观测性较弱你不知道什么时候链路断开了排查问题不方便。手动传入Callback Handler则可以在代码里明确看到监控逻辑而且可以针对特定请求做精细化的追踪配置。from langfuse.callback import CallbackHandler from langfuse import Langfuse langfuse_handler CallbackHandler( public_keypk-..., secret_keysk-..., hosthttp://你的服务器IP:3000 )3. 后端核心代码实现3.1 FastAPI应用初始化与CORS配置FastAPI后端是整个系统的中枢。我把它分成两个部分一部分处理前端的聊天请求和WebSocket推送另一部分负责把Langchain的事件桥接到WebSocket。先看主入口文件的代码。from fastapi import FastAPI, WebSocket, WebSocketDisconnect from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import HTMLResponse from fastapi.templating import Jinja2Templates import json import uvicorn from ws_manager import ConnectionManager from langchain_app import build_llm_chain from callbacks import MonitorCallbackHandler app FastAPI() app.add_middleware( CORSMiddleware, allow_origins[*], allow_credentialsTrue, allow_methods[*], allow_headers[*], ) templates Jinja2Templates(directorytemplates) manager ConnectionManager()CORS的allow_origins这里我直接用了通配符。生产环境建议收紧只允许你的前端域名访问否则别人也可以直接连你的WebSocket接口会有被刷流量的风险。3.2 构建Langchain对话链路对话链路的核心是把提示词模板、LLM实例和可选的记忆模块串起来。我用的是Langchain的RunnableSequence来组装先把用户输入格式化到提示词模板再传给LLM最后解析输出。代码示例如下from langchain_core.prompts import ChatPromptTemplate from langchain_openai import ChatOpenAI from langchain.schema import StrOutputParser from langchain.memory import ConversationBufferMemory from langchain_core.runnables import RunnablePassthrough def build_llm_chain(): prompt ChatPromptTemplate.from_messages([ (system, 你是一个智能助手请用简洁专业的语言回答用户问题。), (human, {question}), ]) llm ChatOpenAI( modeldeepseek-chat, api_keysk-..., base_urlhttps://api.deepseek.com, temperature0.7, streamingTrue, ) memory ConversationBufferMemory(return_messagesTrue) chain ( RunnablePassthrough.assign( chat_historylambda x: memory.load_memory_variables({})[history] ) | prompt | llm | StrOutputParser() ) return chain, memory这里用RunnablePassthrough把历史记忆注入到提示词中再用ChatPromptTemplate组合system和human的消息。StrOutputParser把模型输出解析成字符串避免返回一串带元数据的对象。3.3 Langfuse手动埋点与自定义Callback完成Langchain链路后我在每次调用时传入两个Callback Handler一个是Langfuse官方提供的用于自动上报trace另一个是我自己写的MonitorCallbackHandler用于把token和耗时事件实时推送到WebSocket。两者的职责要分清楚Langfuse管持久化存储和可视化自定义Callback管实时推送。from langchain_core.callbacks import BaseCallbackHandler import json class MonitorCallbackHandler(BaseCallbackHandler): def __init__(self, ws_manager): self.ws_manager ws_manager async def on_llm_new_token(self, token: str, **kwargs) - None: await self.ws_manager.broadcast(json.dumps({ type: token, data: token }, ensure_asciiFalse)) async def on_llm_end(self, response, **kwargs) - None: # 从llm_output中提取token使用量 llm_output kwargs.get(output, {}) if hasattr(response, llm_output) and response.llm_output: token_usage response.llm_output.get(token_usage, {}) await self.ws_manager.broadcast(json.dumps({ type: meta, data: { prompt_tokens: token_usage.get(prompt_tokens, 0), completion_tokens: token_usage.get(completion_tokens, 0), total_tokens: token_usage.get(total_tokens, 0), } }))自定义Callback最关键的一点是方法命名。Langchain的BaseCallbackHandler定义了on_llm_new_token、on_llm_end、on_chain_start、on_chain_end等钩子方法你只需要重写自己关心的那些就行。我在实际测试中发现on_llm_new_token在流式模式下会逐token触发非常适合做打字机效果而on_llm_end则是在整个LLM调用结束后触发适合上报汇总指标。有一点要特别提醒如果你的链路由多个LLM调用组成每个LLM都会触发一遍这些回调所以最好在Callback里加一个trace_id参数用来区分不同的对话会话。3.4 WebSocket连接管理与并发处理WebSocket的核心是ConnectionManager它负责维护所有活跃连接并提供broadcast方法把消息推给所有客户端。这个模块需要处理一个重要的并发问题当多个浏览器页面同时连接时消息广播要保证线程安全。from fastapi import WebSocket, WebSocketDisconnect from typing import List import asyncio class ConnectionManager: def __init__(self): self.active_connections: List[WebSocket] [] self._lock asyncio.Lock() async def connect(self, websocket: WebSocket): await websocket.accept() async with self._lock: self.active_connections.append(websocket) def disconnect(self, websocket: WebSocket): async with self._lock: if websocket in self.active_connections: self.active_connections.remove(websocket) async def broadcast(self, message: str): async with self._lock: connections self.active_connections.copy() for connection in connections: try: await connection.send_text(message) except Exception: await self.disconnect(connection)锁的使用非常关键。如果不用锁多个协程同时对active_connections进行增删操作会出现列表竞争问题轻则消息丢失重则导致整个WebSocket服务崩溃。这里的copy操作是为了避免在遍历连接列表时修改列表本身。FastAPI的WebSocket端点定义如下app.websocket(/ws/monitor) async def websocket_endpoint(websocket: WebSocket): await manager.connect(websocket) try: while True: data await websocket.receive_text() # 收到前端发送的聊天消息 if data.startswith(chat:): question data[5:] await handle_chat(question) else: # 其他控制消息如ping心跳 await websocket.send_text(json.dumps({type: pong})) except WebSocketDisconnect: manager.disconnect(websocket)handle_chat函数会调用前面构建的Langchain链路并把MonitorCallbackHandler传给invoke方法。async def handle_chat(question: str): chain, memory build_llm_chain() langfuse_handler get_langfuse_handler() monitor_handler MonitorCallbackHandler(manager) async def on_chain_start(): await manager.broadcast(json.dumps({ type: status, data: thinking })) await on_chain_start() result await chain.ainvoke( {question: question}, config{callbacks: [langfuse_handler, monitor_handler]} ) memory.chat_memory.add_user_message(question) memory.chat_memory.add_ai_message(result) await manager.broadcast(json.dumps({ type: done, data: result }, ensure_asciiFalse))3.5 流式输出与消息协议设计前端页面的实时渲染依赖于一个清晰的消息协议。我给WebSocket定义了三种消息类型token表示流式输出的token片段meta表示token消耗等元数据done表示整轮回答结束。前端根据type字段做不同的渲染逻辑。token消息可以做成打字机效果meta消息用来更新仪表盘上的累计token数和成本统计done消息用来结束当前的loading状态并展示最终完整答案。协议设计的核心原则是消息要足够小避免大JSON对象阻塞WebSocket通道同时要具备自描述性前端收到消息后不需要额外的上下文就能渲染。4. 前端仪表盘与实时交互4.1 轻量前端方案选择前端我没有引入React或Vue这类重型框架而是用了一个HTML文件加原生JavaScript。原因很简单这个页面只需要做几件事——连WebSocket、接收消息、更新DOM、画几张小图表。用框架反而增加构建复杂度还得配Node环境。Jinja2模板直接返回HTML页面原生JS搞定一切。!DOCTYPE html html langzh head meta charsetUTF-8 titleAI对话监控仪表盘/title script srchttps://cdn.jsdelivr.net/npm/echarts5/dist/echarts.min.js/script /head body div classlayout div classsidebar h3实时对话流/h3 div idchat-log/div /div div classmain-panel div classmetrics div idtoken-chart styleheight: 200px;/div div idlatency-chart styleheight: 200px;/div /div div classinput-area textarea idquestion-input placeholder输入测试问题.../textarea button idsend-btn发送/button /div /div /div /body4.2 WebSocket客户端与断线重连前端WebSocket客户端的核心是断线重连逻辑。我这里采用了指数退避策略断线后先等1秒重连失败后再等2秒、4秒、8秒最多不超过30秒。这样既能在服务端重启后自动恢复连接又不会在服务端故障时疯狂重连打挂服务器。let ws null; let retryCount 0; function connectWebSocket() { const protocol location.protocol https: ? wss : ws; ws new WebSocket(${protocol}://${location.host}/ws/monitor); ws.onopen function() { retryCount 0; console.log(WebSocket connected); }; ws.onmessage function(event) { const msg JSON.parse(event.data); handleMessage(msg); }; ws.onclose function(event) { if (event.code ! 1000) { // 非正常关闭触发重连 const delay Math.min(1000 * Math.pow(2, retryCount), 30000); retryCount; setTimeout(connectWebSocket, delay); } }; } connectWebSocket();这里我特别关注了1006这个错误码。在WebSocket协议里1006表示连接异常关闭通常意味着网络中断或者服务端没有发送关闭帧就断开了。浏览器对这个错误码的处理方式是不触发onerror只在onclose里暴露出来。我在实际运行中发现FastAPI服务端用uvicorn单进程跑的时候代码重新加载或者部署重启时客户端就会收到1006。所以重连逻辑必须覆盖1006场景否则页面会一直卡在断开状态。4.3 指标可视化与效果验证指标可视化我用了ECharts。它会维护一个数据数组每当收到meta消息时就把最新数据push进去然后更新图表。token消耗曲线通常呈现阶梯上升的趋势因为每次对话请求都会增加新的token延迟曲线则可以看出每轮请求的响应时间波动。前端页面整体效果是上面两个图表实时滚动更新下方一个对话流区域逐token展示模型输出旁边显示累计的token消耗和当前轮次的耗时。对于演示和调试来说这个效果已经完全够用了。5. 常见问题与排查技巧5.1 WebSocket 1006断连与重连失效这是我在整个项目中踩得最深的坑。现象是浏览器控制台打印onclosecode为1006reason为空字符串reconnect标志为true。表面上看是网络层断连实际上很多时候是服务端主动断了连接但没有发close帧。排查步骤我整理了一个优先级顺序检查服务端日志看是否有WebSocketDisconnect异常或者uvicorn报错的堆栈信息。检查nginx或其他反向代理层的proxy_read_timeout配置如果设置得太短长时间没有消息传输的连接会被强制断开。在连接管理器里加ping/pong心跳。浏览器WebSocket不支持主动ping但可以在前端定时发ping消息服务端收到后回pong这样连接会一直保持活跃。我最后的解决方案是双管齐下前端加心跳每30秒发送一次ping服务端在收发消息时打日志确认连接状态。实测下来1006断连的次数大幅减少即使偶尔断掉也能迅速重连。5.2 Langfuse看不到Trace数据刚开始集成Langfuse时我遇到了一个问题链路的代码执行了Langfuse后台却看不到任何Trace。排查后发现是回调没有正确传入。Langfuse的Trace数据是通过Callback机制采集的如果你在调用chain时忘了在config里传callbacks参数它就不会上报。另外有一个容易忽略的地方CallbackHandler的初始化时机。如果在创建chain时就把handler传入而运行chain时又没用这个chain只调用了一个独立的llm方法那一样不会上报。正确做法是确保每个需要被监控的调用都显式传入了handler或者干脆创建handler后通过langchain的set_verbose全局设置。Langfuse异步模式下还有一个坑如果你用FastAPI的异步接口调用chain而CallbackHandler不是异步安全的可能会丢失部分trace。解决办法是使用Langfuse提供的AsyncCallbackHandler或者在事件循环中正确调度任务。5.3 FastAPI热更新不生效开发时我一度以为FastAPI不支持热更新改了代码不重启服务就没反应。后来才发现是uvicorn启动命令的问题。如果你用uvicorn.run(app, host, port)这种方式启动默认不会开启reload模式。需要加一个参数uvicorn.run(app, host0.0.0.0, port8000, reloadTrue)。但要注意reloadTrue会启动一个额外的watchdog进程来监听文件变化如果你的代码里有全局变量初始化重载后这些变量会被重新执行一遍。对于WebSocket连接来说重载会导致现有连接全部断开前端重连逻辑此时一定要到位否则就要手动刷新页面。5.4 异步任务挂起与队列堆积设备上报的对话请求多了以后我发现一个严重问题WebSocket广播任务会阻塞模型调用。原因是自定义Callback的on_llm_new_token方法是async的它内部调用了broadcast而broadcast又需要遍历所有连接并发送消息。当连接数多或者网络慢时这个await会拖慢整个链路的执行。解决思路是把广播操作放到独立的任务队列中让Callback不直接等待发送完成。我用的是asyncio.create_task来异步执行广播避免阻塞主流程。但这样做带来另一个问题如果广播任务堆积太多内存会涨。所以最后还是做了个妥协——在同一时刻只保留最近的N条广播消息保证实时性的同时限制内存占用。5.5 常用排查速查表症状可能原因快速检查方式Langfuse后台无TraceCallback未传入或Handler初始化失败检查日志是否输出Langfuse错误信息WebSocket频繁1006反向代理超时或心跳缺失添加心跳调整proxy_read_timeout前端token渲染卡顿广播任务阻塞了模型调用改用create_task异步广播热更新不生效uvicorn缺少reload参数启动命令加--reload模型首字延迟过高DeepSeek API负载高或网络问题在Langfuse看generation耗时token统计不准确流式模式下token统计需从llm_output获取检查on_llm_end回调的response结构6. 项目体验与进阶扩展这个项目跑通之后我最大的感受是大模型应用的可观测性一定要在项目初期就设计进去而不是等上线出问题了再补。Langfuse的Trace可视化能力确实很强配合WebSocket的实时推送调试体验比我之前看日志的方式好太多了。以前排查一个问题要拉半天日志现在直接在仪表盘上看到整条链路的输入输出、耗时和token消耗问题定位基本可以做到分钟级。如果你想在这个基础上继续扩展我有几个建议可以参考。一是接入Langfuse的Prompt管理功能把项目的提示词都在Langfuse后台维护代码里通过SDK拉取这样调整提示词不用改代码重启服务。二是把监控数据接入告警系统比如在单次请求token消耗超过阈值或者模型响应超时的时候自动通过邮件或企业微信机器人推送告警这样线上问题能第一时间感知。三是做一个多项目隔离的版本给不同业务线分配不同的Langfuse项目监控面板按业务维度筛选数据。我个人的体会是技术栈本身并不复杂难的是把每个组件之间的边界理清楚。Langchain负责编排模型调用Langfuse负责记录和可视化FastAPI负责桥接和实时转发DeepSeek负责提供推理能力WebSocket负责最后一百米的推送每个组件各司其职整条链路的可维护性就会非常高。希望这篇文章能帮你避开我踩过的那些坑顺利构建出自己的AI对话监控系统。
返回列表