
1. 从 Gemini Live Avatar 说起实时 AI 的真实门槛在哪里Gemini Live Avatar 这类产品刚出来的时候很多人的第一反应是“不就是把聊天窗口换成了数字人视频吗”。我一开始也这么想直到自己动手把一套文本工作流往实时交互方向改造才发现事情远没有这么简单。把聊天变成视频本质上只是换了一层皮真正难的是让 AI 在说话的同时还能调用工具、查询数据、更新状态并且这一切要在几百毫秒内完成用户还不能感觉到明显的卡顿。这里面的核心矛盾在于文本工作流是同步的、线性的而实时交互是异步的、并发的。你在文本框里输入一句话等三秒钟拿到回复这很正常。但在实时视频对话里用户说完一句话如果 AI 停顿超过一秒半体验就会断崖式下跌。更麻烦的是AI 在回答“今天天气怎么样”的时候背后可能需要调用天气接口、查询用户位置、格式化输出这些操作如果串行执行延迟会直接叠加到用户的等待时间里。RelayRouter 在这个场景里的位置就是解决“实时交互层”和“后端工具层”之间的调度问题。它不是一个具体的框架而是一种架构模式把前端的实时连接WebSocket 或 SSE与后端的异步工具调用解耦让 AI 的语音/视频输出可以先行工具调用的结果后续再注入。这样用户听到的是一句“我帮你查一下”而不是干等三秒的沉默。适合读这篇内容的人大概分三类一是正在做 AI 对话产品、想从纯文本升级到实时交互的开发者二是后端工程师需要理解 WebSocket 和异步任务队列怎么配合三是对实时 AI 架构感兴趣、想搞清楚“为什么我的 WebSocket 推送总是慢半拍”的同行。我会从架构设计、核心细节、实操实现、问题排查四个层面展开尽量把踩过的坑和验证过的方案都写清楚。2. 整体架构设计为什么 RelayRouter 不是简单的消息转发2.1 实时 AI 工作流的三个层次先把结构拆开看。一套完整的实时 AI 工作流我习惯把它分成三层交互层负责和用户直接打交道包括语音识别、视频渲染、文本输入。这一层的关键指标是延迟通常要求端到端控制在 800ms 以内。调度层也就是 RelayRouter 所在的位置。它接收交互层的请求判断哪些需要调用外部工具哪些可以直接由模型生成然后把任务分发给后端。工具层实际执行查询、计算、写入的操作比如查数据库、调第三方 API、读写文件。很多团队一开始会把这三层揉在一起前端直接连后端接口后端同步调用工具结果就是任何一个环节变慢整个链路都卡住。RelayRouter 的思路是把调度层独立出来让它成为一个“异步中转站”。2.2 为什么选择 WebSocket 而不是纯 HTTP这里要解释一个关键选型。实时 AI 场景下WebSocket 几乎是默认选项但很多人说不清楚为什么不用 HTTP 轮询或者 SSE。HTTP 轮询的问题在于每次请求都要重新建立连接、携带头部信息开销大且延迟不可控。SSE 是单向的服务端可以推但客户端不能通过同一个连接发消息实时对话需要双向通道。WebSocket 建立一次连接后双方可以随时互发消息头部开销极小适合高频小消息的场景。但 WebSocket 也不是没有代价。它是有状态的服务端需要维护每个连接的状态这在水平扩展时会带来麻烦。我的做法是WebSocket 只负责传输状态存在 Redis 里。这样即使连接断开重连到另一台服务器也能通过 session id 恢复上下文。2.3 RelayRouter 的核心职责拆解RelayRouter 具体做什么我把它归纳为四个动作消息分类收到用户输入后先判断这是纯聊天、工具调用还是混合请求。纯聊天直接走模型生成工具调用进入异步队列。任务分发把工具调用请求拆成独立任务带上 trace id投递到消息队列我用的是 Redis Stream轻量且够用。结果聚合工具执行完成后结果回到 RelayRouter由它决定是立即推送给前端还是等模型生成完再一起推送。超时兜底如果某个工具调用超过预设时间比如 2 秒RelayRouter 会先给前端发一个“正在处理”的信号避免用户以为卡死。这个设计的优势在于模型生成和工具调用可以并行。比如用户问“帮我总结一下昨天的会议记录并生成待办”模型可以先把“好的我来处理”这句话生成出来同时工具层去读取会议记录等结果回来后再生成具体内容。用户感知到的延迟就从“读取生成”变成了“max(读取, 生成)”。3. 核心细节解析WebSocket 连接管理与异步工具调用3.1 WebSocket 心跳机制的正确实现方式WebSocket 连接如果不维护很容易被中间层断开。我试过最粗暴的方式是前端每 30 秒发一个 ping后端回 pong。但实测下来这种方式在移动网络下经常误判——用户切到后台再回来连接其实已经断了但前端还在傻等。更稳的做法是双向心跳加超时检测// 前端心跳实现 let heartbeatInterval; let lastPongTime Date.now(); function startHeartbeat(ws) { heartbeatInterval setInterval(() { if (Date.now() - lastPongTime 45000) { console.log(心跳超时主动重连); ws.close(); return; } ws.send(JSON.stringify({ type: ping, timestamp: Date.now() })); }, 15000); } ws.onmessage (event) { const data JSON.parse(event.data); if (data.type pong) { lastPongTime Date.now(); return; } // 处理其他消息 };后端对应地要在收到 ping 时立即回 pong并且记录每个连接的最后活跃时间。如果超过 60 秒没有收到任何消息服务端主动关闭连接释放资源。注意心跳间隔不要设得太短15 到 30 秒是比较合理的范围。太短会增加无谓的网络开销太长则断线检测不及时。3.2 异步工具调用的任务队列设计工具调用的异步化是 RelayRouter 的核心。我用 Redis Stream 做任务队列结构大概是这样的# 任务投递示例Python Redis import redis import json import uuid r redis.Redis(hostlocalhost, port6379, decode_responsesTrue) def dispatch_tool_call(session_id, tool_name, params): task_id str(uuid.uuid4()) task { task_id: task_id, session_id: session_id, tool_name: tool_name, params: json.dumps(params), status: pending, created_at: time.time() } r.xadd(tool_tasks, task) return task_id消费者端用消费者组来保证任务不被重复执行def consume_tasks(): group_name tool_workers try: r.xgroup_create(tool_tasks, group_name, id0, mkstreamTrue) except redis.exceptions.ResponseError: pass # 组已存在 while True: messages r.xreadgroup( group_name, worker-1, {tool_tasks: }, count1, block5000 ) for stream, msgs in messages: for msg_id, task in msgs: result execute_tool(task) push_result_to_session(task[session_id], result) r.xack(tool_tasks, group_name, msg_id)这里的关键点是任务执行结果不直接返回给调用方而是推送到 session 对应的通道。RelayRouter 监听这个通道拿到结果后再决定怎么合并到 AI 的输出流里。3.3 文本工作流中的消息合并策略工具调用的结果回来之后怎么和模型的文本输出合并我试过三种策略策略做法适用场景缺点等待合并等所有工具结果回来再生成结果依赖性强延迟高流式插入模型先输出工具结果回来就插入实时对话可能打断语义分段生成模型生成一段工具结果注入后再生成下一段复杂任务实现复杂我目前用的是分段生成。具体做法是模型先输出一个“过渡句”比如“我查到了以下信息”然后暂停生成等工具结果注入后再继续生成具体内容。这样用户听到的是连贯的语音而不是突然跳出一段数据。实现上RelayRouter 维护一个 per-session 的状态机class SessionState: def __init__(self): self.pending_tools set() self.buffer [] self.generation_paused False def on_tool_result(self, task_id, result): self.pending_tools.discard(task_id) self.buffer.append(result) if not self.pending_tools and self.generation_paused: self.resume_generation()这个状态机要处理好并发多个工具同时返回时不能重复触发 resume工具超时时要能强制 resume 并给出降级回复。4. 实操过程从零搭建一套可用的实时文本工作流4.1 环境准备与依赖选型我用的技术栈是前端 React 原生 WebSocket API后端 Python Django Channels处理 WebSocket任务队列 Redis Stream模型调用走异步 HTTP 客户端 httpx。Django Channels 的好处是它把 WebSocket 的消费者模型和 Django 的 ORM 结合得比较好不用自己造轮子。安装pip install channels channels-redis httpx redissettings.py 里配置INSTALLED_APPS [ daphne, channels, # ... ] ASGI_APPLICATION myproject.asgi.application CHANNEL_LAYERS { default: { BACKEND: channels_redis.core.RedisChannelLayer, CONFIG: { hosts: [(127.0.0.1, 6379)], }, }, }提示Daphne 要放在 INSTALLED_APPS 的最前面否则 runserver 不会走 ASGI。4.2 WebSocket 消费者与 RelayRouter 的对接Django Channels 的消费者写法import json from channels.generic.websocket import AsyncWebsocketConsumer from channels.db import database_sync_to_async class ChatConsumer(AsyncWebsocketConsumer): async def connect(self): self.session_id self.scope[url_route][kwargs][session_id] self.group_name fchat_{self.session_id} await self.channel_layer.group_add(self.group_name, self.channel_name) await self.accept() await self.start_heartbeat_check() async def receive(self, text_data): data json.loads(text_data) if data[type] ping: await self.send(json.dumps({type: pong})) return if data[type] user_message: await self.handle_user_message(data[content]) async def handle_user_message(self, content): # 先给前端一个确认信号 await self.send(json.dumps({ type: ack, message: 正在处理... })) # 投递到 RelayRouter await self.channel_layer.send( relay_router, {type: route_message, session_id: self.session_id, content: content} ) async def chat_message(self, event): await self.send(json.dumps(event[data]))RelayRouter 本身是一个独立的消费者监听relay_router通道class RelayRouterConsumer(AsyncWebsocketConsumer): async def route_message(self, event): session_id event[session_id] content event[content] # 判断是否需要工具调用 tool_calls await self.classify_intent(content) if tool_calls: for call in tool_calls: task_id dispatch_tool_call(session_id, call[name], call[params]) # 通知前端有工具在跑 await self.channel_layer.group_send( fchat_{session_id}, {type: chat_message, data: {type: tool_pending, task_id: task_id}} ) # 同时启动模型生成 await self.start_generation(session_id, content)4.3 工具执行结果的回推与前端渲染工具执行完成后结果通过 Redis 的 pub/sub 回到 Django Channels# 工具执行端 def push_result_to_session(session_id, result): channel_layer get_channel_layer() async_to_sync(channel_layer.group_send)( fchat_{session_id}, {type: chat_message, data: {type: tool_result, result: result}} )前端收到tool_result后不直接渲染而是交给一个缓冲区等模型生成到合适的断点时再插入。这个断点判断我用的是标点符号检测模型输出到句号或问号时检查缓冲区是否有待插入的结果。let resultBuffer []; let pendingInsert false; ws.onmessage (event) { const data JSON.parse(event.data); if (data.type tool_result) { resultBuffer.push(data.result); pendingInsert true; return; } if (data.type text_chunk) { appendText(data.content); if (pendingInsert /[。.!?]/.test(data.content)) { const result resultBuffer.shift(); appendText(formatResult(result)); pendingInsert resultBuffer.length 0; } } };4.4 参数计算超时阈值与重试次数怎么定超时阈值不是拍脑袋定的。我的方法是先统计工具调用的 P95 延迟然后取 P95 的 1.5 倍作为超时时间。比如天气接口 P95 是 800ms那超时设 1200ms。超过这个时间RelayRouter 就发一个“还在查”的信号让模型生成一句过渡话术。重试次数我一般设 1 次且只在网络类错误时重试。业务类错误比如参数不对重试没有意义。重试的间隔用指数退避第一次等 200ms第二次等 500ms。注意重试要在任务队列层面做不要在 WebSocket 连接里做。连接断了重试会丢上下文。5. 常见问题与排查技巧实录5.1 WebSocket 连接成功但收不到消息这是最常见的问题。排查顺序我一般是检查 group 是否加入成功在 connect 里打印 channel_name 和 group_name确认没有异常。检查 channel layer 配置Redis 地址对不对能不能 ping 通。检查消息类型是否匹配Channels 的 group_send 里的 type 要和消费者方法名对应比如chat_message对应async def chat_message。检查前端是否在正确的 session 里多标签页测试时session_id 可能串了。我踩过最坑的一次是 Redis 的 db 选错了channel layer 连的是 db 0任务队列连的是 db 1结果消息发到了另一个库前端一直等不到。5.2 工具调用结果乱序或重复乱序通常是因为多个工具并发执行结果回来的顺序不确定。解决办法是在前端按 task_id 去重并且只插入一次。重复则可能是消费者组没有正确 ack导致消息被重新投递。排查命令# 查看 Redis Stream 的 pending 消息 redis-cli XPENDING tool_tasks tool_workers如果有大量 pending说明消费者处理完没有 ack或者处理时间过长被判定超时。5.3 心跳正常但消息延迟高心跳正常只说明连接没断不代表消息通道畅通。延迟高通常是后端处理慢。我的排查方法是加 trace id从前端发送到后端接收、任务投递、工具执行、结果回推每个环节打时间戳最后算各段耗时。常见瓶颈点环节典型耗时优化方向前端到后端20-50ms检查网络用 CDN 加速意图分类100-300ms用小模型或规则引擎工具执行200-2000ms加缓存异步化结果回推10-30ms检查 channel layer 负载5.4 断线重连后的状态恢复WebSocket 断线后前端要能自动重连并且恢复未完成的任务。我的做法是前端保存最后收到的消息序号seq。重连时带上 seq后端从 Redis 里取 seq 之后的消息补发。如果任务还在执行中后端返回任务状态前端继续等待。function reconnect() { const ws new WebSocket(wss://example.com/ws/chat/${sessionId}?seq${lastSeq}); ws.onopen () { console.log(重连成功请求补发消息); ws.send(JSON.stringify({ type: resume, seq: lastSeq })); }; return ws; }提示seq 用单调递增的整数存在内存里就行不需要持久化。重连后如果 seq 对不上就让后端全量推一次当前会话状态。5.5 常见问题速查表现象可能原因解决方式连接建立后立即断开认证失败或 URL 路由错误检查 token 和 routing 配置消息发送成功但无响应消费者方法名不匹配核对 type 和方法名工具结果重复插入未按 task_id 去重前端加 Set 去重心跳超时误判移动网络切换放宽超时阈值到 60s多标签页消息串扰session_id 复用每个标签页独立 session任务队列积压消费者处理慢增加消费者实例6. 从文本工作流到实时 AI 的扩展思路RelayRouter 这套模式不只适用于聊天场景。我后来把它用在了文档协作工具里用户编辑文档时后台异步跑语法检查和内容建议结果通过 WebSocket 推送到编辑器但不打断用户的输入。核心逻辑是一样的——把耗时的工具调用从主交互链路里剥离出去用异步通道回推结果。另一个扩展方向是多人协作。多个用户的 WebSocket 连接加入同一个 group任何一个人的操作都广播给其他人。这时候 RelayRouter 要处理冲突合并我用的是操作转换OT的简化版每个操作带时间戳和用户 id服务端按时间排序后广播。如果你正在做类似的东西我的建议是先把单用户的异步工具调用跑通再考虑多用户和复杂状态同步。一上来就搞大而全的架构很容易在调试阶段迷失方向。先把 WebSocket 心跳、任务队列、结果回推这三个环节调稳后面的扩展都是水到渠成的事。最后分享一个我实际使用中总结的小技巧在开发阶段把 RelayRouter 的每一步都打上日志并且用一个 debug 通道把日志实时推到前端。这样你在浏览器控制台就能看到消息从进入到返回的完整路径排查问题比翻服务端日志快得多。等稳定了再把 debug 通道关掉不影响性能。