云计算百科
云计算领域专业知识百科平台

从零构建 AI 对话监控仪表盘:Langfuse + Langchain + DeepSeek + FastAPI + WebSocket 实时架构实践

当你的 AI 应用开始有真实用户时,"模型回复了什么"已经不够了——你需要知道谁在用、用了多少 Token、花了多少钱、响应有多快。本文通过一个完整的 Demo 项目,展示如何用 FastAPI + DeepSeek + WebSocket 构建一个带实时监控仪表盘的 AI 对话服务。


一、这个 Demo 做了什么

项目提供两个服务,跑在同一个 FastAPI 应用里:

服务

面向

功能

AI 对话服务

终端用户

对接 DeepSeek 大模型,支持多轮对话、会话管理

实时仪表盘

运营/开发

实时查看用户活跃度、Token 消耗、费用统计、响应延迟

两个服务的数据流是串联的:用户每发一条消息,仪表盘无需刷新就能看到新的活动记录和统计变化。


二、快速启动

2.1 环境要求

  • Python 3.11+
  • DeepSeek API Key(在 platform.deepseek.com 获取)

2.2 安装步骤

# 1. 进入项目目录
cd D:\\Dashboard

# 2. 创建虚拟环境
python -m venv venv
venv\\Scripts\\activate # Windows
# source venv/bin/activate # macOS/Linux

# 3. 安装依赖
pip install -r requirements.txt

# 4. 配置环境变量
copy .env.example .env

编辑 .env 文件,填入你的 DeepSeek API Key:

DEEPSEEK_API_KEY=sk-your-actual-key-here
DEEPSEEK_MODEL=deepseek-chat

2.3 启动服务

python run.py

启动后访问以下地址:

地址

说明

http://localhost:8000/chat

用户对话界面

http://localhost:8000/dashboard

实时监控仪表盘

High Performance Web Crawler API – Swagger UI

API 交互文档(Swagger UI)

2.4 使用流程

  • 打开 /chat 页面,输入用户名(如"张三")
  • 输入消息,如"帮我写一首关于秋天的诗"
  • DeepSeek 返回回复,界面底部显示本次调用的 Token 数和费用
  • 打开 /dashboard 页面(可以开在新标签页)
  • 在对话页发一条新消息 → 仪表盘的"实时活动"区域立即出现这条记录
  • 统计卡片、Token 趋势图、费用趋势图也会同步刷新

  • 三、项目结构

    整体采用分层架构:路由层(routers)负责 HTTP 入口,服务层(services)封装业务逻辑,数据层(models + database)管理持久化,WebSocket 层(ws)处理实时通信。各层之间通过依赖注入解耦。


    四、关键技术环节详解

    4.1 DeepSeek API 集成:复用 OpenAI SDK

    DeepSeek 的 API 完全兼容 OpenAI 接口格式,这意味着不需要安装额外的 SDK——直接用 openai 官方库,把 base_url 指向 DeepSeek 即可。

    # app/services/deepseek.py

    from openai import AsyncOpenAI

    class DeepSeekService:
    def __init__(self):
    self.client = AsyncOpenAI(
    api_key=settings.deepseek_api_key,
    base_url="https://api.deepseek.com", # 关键:指向 DeepSeek
    )
    self.model = "deepseek-chat"

    async def chat(self, messages: list[dict]) -> DeepSeekResult:
    start = time.monotonic()

    response = await self.client.chat.completions.create(
    model=self.model,
    messages=messages,
    )

    latency_ms = int((time.monotonic() – start) * 1000)

    return DeepSeekResult(
    content=response.choices[0].message.content,
    model=self.model,
    tokens_input=response.usage.prompt_tokens, # 输入 Token
    tokens_output=response.usage.completion_tokens, # 输出 Token
    latency_ms=latency_ms,
    )

    关键点:

    • 使用 AsyncOpenAI 而非同步客户端,确保 API 调用不阻塞事件循环
    • time.monotonic() 而非 time.time() 测量延迟——前者不受系统时钟调整影响
    • 从 response.usage 直接提取 Token 用量,无需自己实现分词计数
    • 返回值用 dataclass 封装,类型清晰,便于下游消费

    4.2 异步数据库:SQLAlchemy 2.0 + aiosqlite

    整个数据层基于异步实现,避免 I/O 阻塞 FastAPI 的事件循环。

    引擎与会话工厂:

    # app/database.py

    from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker, AsyncSession

    engine = create_async_engine("sqlite+aiosqlite:///./dashboard.db", echo=False)
    async_session = async_sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)

    async def get_db() -> AsyncSession:
    """FastAPI 依赖注入入口"""
    async with async_session() as session:
    yield session

    关键设计:

    • expire_on_commit=False:commit 之后 ORM 对象仍然可用(默认为 True 会导致 commit 后访问属性触发额外查询)
    • get_db() 是一个异步生成器,通过 FastAPI 的 Depends 注入到每个路由函数中,请求结束自动关闭会话
    • 使用 SQLite + aiosqlite 作为零配置存储,生产环境可无缝切换为 PostgreSQL(改连接字符串即可)

    数据模型:

    # app/models.py(精简)

    class User(Base):
    id: Mapped[int] # 自增主键
    user_id: Mapped[str] # 业务用户ID(如 "张三")
    name: Mapped[str]
    created_at / last_active_at

    class ChatSession(Base):
    session_id: Mapped[str] # 业务会话ID(如 "sess_a1b2c3…")
    user_id → ForeignKey("users.user_id")
    title: Mapped[str] # 取首条消息前30字作为标题

    class Message(Base):
    session_id → ForeignKey("sessions.session_id")
    role: Mapped[str] # "user" 或 "assistant"
    content: Mapped[str]
    model: Mapped[str] # 记录是哪个模型回复的
    tokens_input / tokens_output # Token 用量
    cost: Mapped[float] # 本次调用费用(元)
    latency_ms: Mapped[int] # 响应延迟(毫秒)
    created_at: Mapped[datetime]

    三层模型形成 User → ChatSession → Message 的归属链。关键索引 ix_messages_session_time 覆盖了"按会话查消息、按时间排序"的高频查询路径。

    4.3 用量追踪编排:track_chat 方法

    tracker.track_chat() 是整个系统的核心编排方法,一次调用串联了"建用户 → 建会话 → 存用户消息 → 调 DeepSeek → 存 AI 回复 → 推送事件"的完整流程:

    # app/services/tracker.py(精简)

    class UsageTracker:
    async def track_chat(self, db, user_id, user_name, session_id,
    user_message, history) -> dict:
    # 1. 确保用户存在(不存在则创建,已存在则更新活跃时间)
    await self.ensure_user(db, user_id, user_name)

    # 2. 确保会话存在(session_id 为空则自动创建)
    session = await self.ensure_session(db, session_id, user_id, title)

    # 3. 先存用户消息
    await self.record_message(db, session.session_id, "user", user_message)

    # 4. 调用 DeepSeek(拼接历史上下文 + 当前消息)
    messages_for_api = history + [{"role": "user", "content": user_message}]
    result = await deepseek_service.chat(messages_for_api)
    cost = deepseek_service.calculate_cost(result.tokens_input,
    result.tokens_output)

    # 5. 存 AI 回复(含 Token/费用/延迟)
    await self.record_message(db, session.session_id, "assistant",
    result.content, model=result.model,
    tokens_input=result.tokens_input,
    tokens_output=result.tokens_output,
    cost=cost, latency_ms=result.latency_ms)
    await db.commit()

    # 6. WebSocket 广播实时事件
    await ws_manager.broadcast("activity", {
    "user_name": user_name,
    "message_preview": user_message[:80],
    "tokens_input": result.tokens_input,
    "tokens_output": result.tokens_output,
    "cost": cost,
    "latency_ms": result.latency_ms,

    })

    return {"reply": result.content, "cost": cost, …}

    关键设计:

    • 先记录后提交:用户消息和 AI 回复在同一个事务中写入,要么都成功,要么都回滚
    • commit 在 broadcast 之前:确保推送到仪表盘的数据已经持久化,不会出现"仪表盘显示了但数据库没有"的不一致
    • broadcast 在 commit 之后:如果广播失败(如 WebSocket 断开),不影响数据持久化
    • 费用计算:基于 DeepSeek 官方定价(输入 1 元/百万 Token,输出 2 元/百万 Token),在 config.py 中可配置

    4.4 WebSocket 实时推送架构

    这是仪表盘"实时"能力的核心。整个推送链路分为三层:

    第一层:连接管理器(ws/manager.py)

    class WebSocketManager:
    def __init__(self):
    self._connections: list[WebSocket] = [] # 所有活跃连接

    async def connect(self, ws: WebSocket):
    await ws.accept()
    self._connections.append(ws)

    def disconnect(self, ws: WebSocket):
    if ws in self._connections:
    self._connections.remove(ws)

    async def broadcast(self, event_type: str, data: dict):
    if not self._connections:
    return # 没有连接时直接跳过,零开销

    payload = json.dumps({"type": event_type, "data": data}, ensure_ascii=False)
    dead = []
    for ws in self._connections:
    try:
    await ws.send_text(payload)
    except Exception:
    dead.append(ws) # 收集发送失败的连接
    for ws in dead:
    self.disconnect(ws) # 清理死连接

    关键点:

    • broadcast 在没有连接时直接返回,不影响对话服务性能
    • 发送失败的连接被收集到 dead 列表,广播结束后统一清理,避免遍历过程中修改列表
    • ensure_ascii=False 确保中文消息正常传输

    第二层:WebSocket 端点(routers/dashboard.py)

    @router.websocket("/ws")
    async def dashboard_ws(ws: WebSocket):
    await ws_manager.connect(ws)
    try:
    while True:
    await ws.receive_text() # 保持连接,等待客户端消息
    except WebSocketDisconnect:
    ws_manager.disconnect(ws)

    这个端点本身不处理客户端发来的消息(仪表盘只需要接收推送),receive_text() 的作用是维持连接存活。当客户端断开时抛出 WebSocketDisconnect 异常,触发清理。

    第三层:前端消费(dashboard.js)

    function connectWebSocket() {
    const protocol = location.protocol === 'https:' ? 'wss:' : 'ws:';
    const ws = new WebSocket(`${protocol}//${location.host}/api/dashboard/ws`);

    ws.onopen = () => {
    document.getElementById('wsStatus').textContent = '已连接';
    };

    ws.onmessage = (event) => {
    const msg = JSON.parse(event.data);
    if (msg.type === 'activity') {
    addActivityItem(msg.data); // 插入到活动流顶部
    loadStats(); // 刷新统计卡片
    loadUsers(); // 刷新用户排行
    }
    };

    ws.onclose = () => {
    document.getElementById('wsStatus').textContent = '已断开';
    setTimeout(connectWebSocket, 5000); // 5秒后自动重连
    };
    }

    前端关键点:

    • 自动选择 ws:// 或 wss:// 协议,适配 HTTP/HTTPS 部署
    • 收到消息后除了插入活动流,还主动拉取最新的统计数据和用户列表(因为 WebSocket 只推送"增量事件",聚合统计需要重新查询)
    • 断线后 5 秒自动重连,保证长时间打开仪表盘时的连接可靠性

    4.5 多轮对话上下文管理

    对话服务支持多轮对话,关键在于如何构建上下文:

    # app/routers/chat.py

    @router.post("")
    async def chat(req: ChatRequest, db: AsyncSession = Depends(get_db)):
    history: list[dict] = []

    if req.session_id:
    # 从数据库加载该会话的历史消息
    result = await db.execute(
    select(Message)
    .where(Message.session_id == req.session_id)
    .order_by(Message.id)
    )
    msgs = result.scalars().all()
    # 只取最近 10 条作为上下文,控制 Token 消耗
    history = [
    {"role": m.role, "content": m.content}
    for m in msgs[-10:]
    ]

    result = await tracker.track_chat(
    db=db, user_id=req.user_id, user_message=req.message,
    session_id=req.session_id, history=history,
    )
    return ChatResponse(**result)

    关键设计:

    • 上下文从数据库持久化加载,而非依赖内存——服务重启后对话不丢失
    • msgs[-10:] 只取最近 10 条消息,避免上下文无限膨胀导致 Token 消耗失控
    • session_id 首次请求时为 None,由 tracker.ensure_session() 自动创建并返回,后续请求携带此 ID 即可延续对话

    4.6 仪表盘聚合查询

    仪表盘的统计接口大量使用 SQL 聚合函数,直接在数据库层面完成计算,避免把全量数据拉到 Python 中处理:

    # 总览统计 —— 一次请求拿到 7 个指标
    @router.get("/stats")
    async def get_stats(db: AsyncSession = Depends(get_db)):
    total_users = await db.scalar(select(func.count(User.id)))
    total_cost = await db.scalar(
    select(func.coalesce(func.sum(Message.cost), 0.0))
    )
    avg_latency = await db.scalar(
    select(func.coalesce(func.avg(Message.latency_ms), 0.0))
    .where(Message.role == "assistant")
    )
    active_24h = await db.scalar(
    select(func.count(func.distinct(ChatSession.user_id)))
    .where(ChatSession.updated_at >= cutoff)
    )

    # 用户用量排行 —— 三表 JOIN + GROUP BY
    stmt = (
    select(
    ChatSession.user_id,
    User.name.label("name"),
    func.count(Message.id).label("message_count"),
    func.sum(Message.tokens_input).label("tokens_input"),
    func.sum(Message.cost).label("total_cost"),
    func.avg(Message.latency_ms).label("avg_latency_ms"),
    )
    .select_from(Message)
    .join(ChatSession, Message.session_id == ChatSession.session_id)
    .join(User, ChatSession.user_id == User.user_id)
    .where(Message.role == "assistant")
    .group_by(ChatSession.user_id, User.name)
    .order_by(desc("total_cost"))
    )
    # 24小时时间线 —— 按小时聚合
    stmt = (
    select(
    func.strftime("%Y-%m-%dT%H:00:00", Message.created_at).label("ts"),
    func.count(Message.id).label("messages"),
    func.sum(Message.tokens_input + Message.tokens_output).label("tokens"),
    func.sum(Message.cost).label("cost"),
    )
    .where(Message.created_at >= cutoff)
    .group_by("ts")
    .order_by("ts")
    )

    关键点:

    • func.coalesce(sum(…), 0) 确保无数据时返回 0 而非 None
    • func.distinct() 用于统计去重指标(如 24h 活跃用户数)
    • func.strftime() 在 SQLite 中做时间截断聚合,按小时分组
    • 延迟平均值只统计 role == "assistant" 的消息,因为用户消息没有 API 调用延迟

    4.7 前端可视化:Chart.js 趋势图

    仪表盘使用 Chart.js 绘制 Token 和费用的时间趋势图:

    async function loadTimeline() {
    const res = await fetch(`${API}/timeline?hours=24`);
    const points = await res.json();

    const labels = points.map(p => {
    const d = new Date(p.timestamp);
    return d.getHours().toString().padStart(2, '0') + ':00';
    });

    tokenChart = new Chart(tokenCtx, {
    type: 'line',
    data: {
    labels: labels,
    datasets: [{
    label: 'Tokens',
    data: points.map(p => p.tokens),
    borderColor: '#4f6ef7',
    backgroundColor: 'rgba(79,110,247,0.08)',
    fill: true,
    tension: 0.3, // 平滑曲线
    pointRadius: 2,
    }]
    },
    options: {
    responsive: true,
    maintainAspectRatio: false,
    plugins: { legend: { display: false } },
    }
    });
    }

    时间线每 60 秒自动刷新一次,配合 WebSocket 的实时推送,形成"实时事件流 + 定时趋势刷新"的双层更新机制。

    4.8 配置管理与定价模型

    使用 pydantic-settings 管理配置,支持环境变量注入:

    # app/config.py

    class Settings(BaseSettings):
    deepseek_api_key: str = os.getenv("DEEPSEEK_API_KEY", "")
    deepseek_base_url: str = "https://api.deepseek.com"
    deepseek_model: str = "deepseek-chat"

    # DeepSeek 定价(元/百万 Token)
    price_input_per_million: float = 1.0 # 输入
    price_output_per_million: float = 2.0 # 输出

    database_url: str = "sqlite+aiosqlite:///./dashboard.db"
    host: str = "0.0.0.0"
    port: int = 8000

    费用计算逻辑:

    def calculate_cost(self, tokens_input: int, tokens_output: int) -> float:
    cost_in = tokens_input / 1_000_000 * settings.price_input_per_million
    cost_out = tokens_output / 1_000_000 * settings.price_output_per_million
    return round(cost_in + cost_out, 6)

    切换模型时只需修改 deepseek_model 和定价参数,无需改代码。


    五、API 接口一览

    对话服务

    方法

    路径

    说明

    POST

    /api/chat

    发送消息,获取 AI 回复

    GET

    /api/chat/sessions

    获取会话列表(支持 user_id 过滤)

    GET

    /api/chat/sessions/{id}

    获取会话详情(含完整消息历史)

    POST /api/chat 请求体:

    {
    "message": "帮我写一首关于秋天的诗",
    "user_id": "zhangsan",
    "user_name": "张三",
    "session_id": null
    }

    响应:

    {
    "session_id": "sess_a1b2c3d4e5f6",
    "reply": "秋风起处叶纷飞…",
    "model": "deepseek-chat",
    "tokens_input": 52,
    "tokens_output": 128,
    "cost": 0.000308,
    "latency_ms": 1840
    }

    仪表盘服务

    方法

    路径

    说明

    GET

    /api/dashboard/stats

    总览统计(7 个指标)

    GET

    /api/dashboard/users

    用户用量排行

    GET

    /api/dashboard/timeline

    24 小时时间线(按小时聚合)

    GET

    /api/dashboard/recent-activity

    最近 20 条活动

    WS

    /api/dashboard/ws

    WebSocket 实时推送端点


    六、数据流全景

    把所有环节串起来,一次完整的对话请求经过以下步骤:

    1. 用户在 /chat 页面输入消息
    2. chat.js 发送 POST /api/chat
    3. chat.py 路由收到请求,从 DB 加载最近 10 条历史消息
    4. 调用 tracker.track_chat():
    a. ensure_user() → 创建或更新用户记录
    b. ensure_session() → 创建或更新会话记录
    c. record_message() → 存储用户消息
    d. deepseek_service.chat() → 异步调用 DeepSeek API
    e. calculate_cost() → 计算费用
    f. record_message() → 存储 AI 回复(含 Token/费用/延迟)
    g. db.commit() → 提交事务
    h. ws_manager.broadcast() → WebSocket 推送事件
    5. 返回 JSON 响应给 chat.js
    6. chat.js 更新对话界面,显示回复和用量
    7. dashboard.js 收到 WebSocket 消息:
    a. 插入活动流顶部
    b. 刷新统计卡片
    c. 刷新用户排行

    整个过程中,DeepSeek API 调用是最慢的一环(通常 1-3 秒),其他环节都是毫秒级。由于全部异步实现,即使有多个用户同时对话,也不会互相阻塞。


    七、扩展方向

    这个 Demo 是一个最小可运行版本,以下是一些实用的扩展方向:

    方向

    做法

    流式输出

    使用 DeepSeek 的 stream=True + FastAPI 的 StreamingResponse,让用户逐字看到回复

    用户认证

    加入 JWT 认证,user_id 从 Token 中提取而非前端传入

    数据库迁移

    引入 Alembic 管理表结构变更,支持平滑升级

    多模型支持

    在 config.py 中配置多个模型,前端可选择,统计按模型分组

    告警机制

    当某用户 Token 消耗超过阈值时,通过 WebSocket 推送告警事件

    数据导出

    增加导出 CSV/Excel 接口,方便离线分析

    部署上线

    用 Gunicorn + Uvicorn worker 部署,数据库切换为 PostgreSQL


    八、总结

    这个 Demo 展示了构建 AI 应用监控体系的几个核心模式:

  • API 兼容复用:DeepSeek 兼容 OpenAI 接口,直接用 openai SDK 即可集成,零额外学习成本
  • 全链路异步:从数据库(aiosqlite)到 HTTP(FastAPI)到 LLM 调用(AsyncOpenAI),全异步实现,高并发友好
  • 追踪嵌入式设计:用量追踪逻辑内嵌在对话流程中(track_chat 方法),业务代码无需感知
  • WebSocket 实时推送:轻量级的连接管理器 + 广播模式,实现仪表盘的实时更新
  • SQL 聚合下沉:统计计算直接用 SQL 聚合函数完成,避免把数据拉到应用层处理
  • 项目的核心价值在于:它把"AI 对话"和"用量监控"这两件事用最小代价串联起来了。如果你正在做一个接入大模型的应用,这个架构可以直接作为起点。

    九、源代码地址

    Dashboard

    赞(0)
    未经允许不得转载:网硕互联帮助中心 » 从零构建 AI 对话监控仪表盘:Langfuse + Langchain + DeepSeek + FastAPI + WebSocket 实时架构实践
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!