首页 / 文章 / 利用Redis流作为每轮对话的会合点,实现可恢复的LLM流式处理。

利用Redis流作为每轮对话的会合点,实现可恢复的LLM流式处理。

通过将代理事件发布到针对单轮对话的Redis流中,使客户端能够在连接中断或工具暂停数分钟后继续运行。

2049 词

我们的目标

智能体回复需要数十秒时间:进行推理、加载技能、调用工具、等待以及处理令牌流。演示版本仅保持一个开放连接,而生产环境会在回复过程中中断,需重新部署网关,并为那些数分钟后才返回结果的工具暂停等待。我们的目标是实现可恢复的实时流式传输,让客户端能够重新连接并继续当前的对话轮次。

发展历程

第一阶段:客户端直接与智能体通信

方式简单但脆弱,任何网络故障都会导致流式传输中断。横向扩展则会导致会话粘性问题或事件丢失。

第二阶段:服务间通过gRPC进行流式传输

内部接口有所改进,但对浏览器客户端而言仍不够便捷,且在容器重启导致的数分钟延迟情况下表现依然不佳。

为何不选择Kafka?

非常适合处理持久性日志;但对于需要每轮会面且数据保留时间较短、且消费者组难以对应“一个浏览器标签页”的场景来说,其重量超过了实际需求。

第3次迭代(有效方案):每轮使用带名称的会面机制

// Agent output
{
  "type": "tool_result",
  "tool": "product_search",
  "data": {
    "items": [...]
  }
}

// Gateway -> TV
{
  "type": "product_carousel",
  "items": [...]
}

// Gateway -> Mobile
{
  "type": "product_list",
  "items": [...]
}
cursor = last_event_id or "0-0"

while True:
    entries = xread({key: cursor}, block=30_000)

    if not entries:          # the only timeout check point
        check_timeouts()
        continue

    for entry_id, event in entries:
        # writes to the socket; not an ack that the client received it
        sse.send(id=entry_id, data=event.payload)
        cursor = entry_id
        if event.type in TERMINAL:
            return
id: 1755600000123-0
data: {"type":"tool_selected","tool":"search"}

id: 1755600000871-0
data: {"type":"response_block","block":{...}}
GET /sessions/{sid}/turns/{tid}/stream
Last-Event-ID: 1755600000871-0
# turn starts: one atomic step (MULTI/EXEC, or a Lua script)
xadd(key, first_event)
expire(key, GENEROUS_TTL)

# producer finishes: bring it in
expire(key, RECONNECT_TTL)

每个用户轮次都会对应一个以turn_id为键的Redis Stream(或流+消费者组模式)。代理会发布令牌/工具相关事件;网关则从客户端上一次的id处继续读取。重新连接时会从该位置继续处理。即使Pod崩溃,流中仍保留足够的历史数据以完成当前轮次。

关键细节

  • 游标管理规则——客户端需确认最后看到的流ID;首次连接后绝不能从0重新开始。
  • 心跳事件——防止在工具等待期间中间节点关闭空闲连接。
  • 状态转换生命周期——明确的turn_started / turn_paused / turn_completed / turn_failed状态标识。
  • TTL设置——在状态转换完成后再让流过期,避免Redis变成无限存储的档案库。
  • 授权机制——turn_id并非授权依据;需将状态转换与已认证的会话绑定。
  • 解决该问题的方案:多分钟暂停处理

    某个工具将任务交给另一个处理流程,后者数分钟后才响应。传统的HTTP流式传输方式已无法使用。借助Redis Streams,代理会发布暂停事件,客户端仍显示“正在处理中”,之后在重新连接时即可继续执行任务,无需重新运行整个流程。

    我们发送给客户端的内容

    输入的事件包括:令牌、tool_start、tool_result摘要(绝不能包含机密信息)、错误以及任务完成状态。请保持数据量较小;将大型数据存储在对象存储中,仅发送引用地址。

    成本、限制与注意事项

    需注意Redis内存使用情况、最大流长度,以及当多个网关同时处理同一轮任务时的数据扩散问题。应限制每位用户的并发处理轮数。在网关部署后,还需对大量重连请求进行压力测试。

    最终状态

    网关是具备恢复功能的读取器,代理则是写入端,而Redis Streams则充当数据交汇点。实时用户体验不会因那些会导致演示架构失效的普通故障而受到影响。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络连接不稳定。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络不稳定所致。

    运营建议:将最新的流标识符存储在会话cookie或客户端内存中,同时在服务器端也保存一份以便后续回放。当用户反馈“程序卡住了”时,客服人员应直接从流数据中还原该对话环节,而无需让用户再次等待长达十分钟。此外还需添加相关仪表板,展示对话恢复率、被放弃的对话次数以及平均暂停时长,以便团队判断是客服处理速度过慢还是网络不稳定所致。