面向模型流式输出的多订阅网关:从 Socket 连接到可恢复事件分发

简介: 本文探讨如何构建高可靠流式模型网关:以Redis Streams持久化事件,通过SSE向多端(浏览器、审计、日志等)分发;支持断线续传、游标恢复、背压隔离与慢客户端容错,避免“广播即送达”的脆弱性,实现可追溯、可重放、可扩展的事件驱动架构。(239字)

把模型输出直接转发给一个前端页面,通常只需要建立一条 HTTP 流。但当同一份结果需要同时送给浏览器、审计服务、日志系统和人工复核端时,简单的“收到一段就广播一段”很快会暴露问题:某个订阅者变慢,是否拖住所有连接?客户端断线后,已经生成的内容能否继续获取?服务重启后,正在进行的任务是否只能从头开始?

这些问题的共同本质不是“如何打印文本”,而是如何管理一组有生命周期的事件流。本文采用一个小型网关作为示例:上游模型接口产生增量事件,网关持久化事件并向多个订阅者提供 SSE 读取接口。SSE 只是浏览器侧的传输形式,可靠性由事件存储、游标和连接管理共同提供。

这里的模型接口可以是企业自建服务、云端兼容接口或其他中转 API。接入 HaerAPI 时,应以其当前公开文档确认请求格式、可用模型、鉴权方式和数据处理条款,再将其作为上游适配器配置进去,而不要把业务代码绑定到某一家服务的私有字段。

核心原理

连接与文件描述符

在 Linux 中,TCP 连接会表现为进程持有的文件描述符。应用层通常通过事件循环监听可读、可写和异常状态。连接本身只负责传输字节,不负责事件是否已经被消费,也不负责断线后的恢复。因此,不能把“网络连接还在”误认为“消息已经可靠送达”。

网关需要把三个对象分开管理:

  1. 事件:具有唯一 ID、任务 ID、类型和数据内容。
  2. 游标:订阅者已经确认读取到的位置。
  3. 连接:当前用于发送事件的 HTTP 或 Socket 通道。

事件先写入可回放的存储,再发送给订阅者。客户端重连时带上上一次收到的事件 ID,服务端从该位置之后继续读取。这样,连接短暂中断不会直接导致业务状态丢失。

为什么选择 Redis Streams

Redis Streams 提供追加式消息记录和消息 ID,适合演示单任务事件回放。生产环境仍需根据保留周期、容量、故障模型和一致性要求评估 Kafka、数据库或其他日志系统。本文不把 Redis Streams 视为所有场景的唯一答案。

事件 ID 通常形如 时间戳-序号。消费者不应自行比较业务时间,而应使用存储系统返回的游标。对于同一个订阅者,读取和更新游标最好具有明确的确认边界:事件成功写入响应缓冲区后可以记为已发送;若业务要求更强的确认语义,则应由客户端回传确认 ID,并接受重复投递后去重。

背压与隔离

SSE 连接的写入速度受客户端网络和浏览器处理速度影响。网关必须限制每个连接的待发送队列。超过上限时,可以断开慢连接并让它通过游标重新追赶,也可以把订阅者转移到独立的持久化消费流程。直接无限扩大内存队列,会把局部网络问题变成整个进程的内存风险。

可执行实现

下面的示例使用 Python、FastAPI 和 Redis。它演示三个接口:创建任务、向任务追加事件、按游标订阅事件。示例把“上游模型调用”留在生产者函数中,便于替换成经过审核的 API 客户端。

安装依赖:

python -m venv .venv
. .venv/bin/activate
pip install fastapi uvicorn redis httpx pydantic-settings

准备环境变量,密钥不写入代码或配置仓库:

export REDIS_URL='redis://localhost:6379/0'
export MODEL_BASE_URL='https://api.example.invalid/v1'
export MODEL_API_KEY='replace-with-runtime-secret'
export MODEL_NAME='your-model-name'

MODEL_BASE_URLMODEL_NAME 只是占位配置。不同服务的兼容程度、流式字段和鉴权头可能不同,应在适配器层处理差异。

服务端示例:

import asyncio
import json
import os
import uuid
from typing import AsyncIterator

import redis.asyncio as redis
from fastapi import FastAPI, Header, HTTPException
from fastapi.responses import StreamingResponse
from pydantic import BaseModel

app = FastAPI()
rdb = redis.from_url(os.environ["REDIS_URL"], decode_responses=True)

class EventIn(BaseModel):
    kind: str
    data: dict

async def append_event(task_id: str, event: EventIn) -> str:
    key = f"task:{task_id}:events"
    return await rdb.xadd(key, {
   
        "kind": event.kind,
        "data": json.dumps(event.data, ensure_ascii=False),
    }, maxlen=10000, approximate=True)

@app.post("/tasks")
async def create_task():
    task_id = str(uuid.uuid4())
    await rdb.hset(f"task:{task_id}:meta", mapping={
   "status": "created"})
    return {
   "task_id": task_id}

@app.post("/tasks/{task_id}/events")
async def publish(task_id: str, event: EventIn):
    if not await rdb.exists(f"task:{task_id}:meta"):
        raise HTTPException(404, "task not found")
    event_id = await append_event(task_id, event)
    return {
   "event_id": event_id}

async def event_stream(task_id: str, last_id: str) -> AsyncIterator[str]:
    key = f"task:{task_id}:events"
    cursor = last_id or "0-0"
    while True:
        rows = await rdb.xread({
   key: cursor}, count=100, block=15000)
        if not rows:
            yield ": keep-alive\\n\\n"
            continue
        for _, items in rows:
            for event_id, fields in items:
                cursor = event_id
                payload = {
   
                    "id": event_id,
                    "kind": fields["kind"],
                    "data": json.loads(fields["data"]),
                }
                yield f"id: {event_id}\\ndata: {json.dumps(payload, ensure_ascii=False)}\\n\\n"

@app.get("/tasks/{task_id}/events")
async def subscribe(task_id: str, last_event_id: str | None = Header(default=None)):
    if not await rdb.exists(f"task:{task_id}:meta"):
        raise HTTPException(404, "task not found")
    return StreamingResponse(
        event_stream(task_id, last_event_id or "0-0"),
        media_type="text/event-stream",
        headers={
   "Cache-Control": "no-cache", "X-Accel-Buffering": "no"},
    )

启动服务:

uvicorn app:app --host 127.0.0.1 --port 8000

创建任务并写入事件:

TASK_ID=$(curl -s -X POST http://127.0.0.1:8000/tasks | jq -r .task_id)
curl -X POST "http://127.0.0.1:8000/tasks/$TASK_ID/events" \
  -H 'Content-Type: application/json' \
  -d '{"kind":"token","data":{"text":"第一段输出"}}'

浏览器或命令行订阅:

curl -N "http://127.0.0.1:8000/tasks/$TASK_ID/events"

浏览器端可以直接使用 EventSource。浏览器会自动携带 Last-Event-ID,但不同代理和客户端实现可能存在差异,因此也可以在应用层显式传递游标:

const source = new EventSource(`/tasks/${
     taskId}/events`);
source.onmessage = (event) => {
   
  const message = JSON.parse(event.data);
  renderDelta(message.data);
};
source.onerror = () => {
   
  // 由浏览器按 EventSource 规则尝试重连;服务端必须支持游标续传
};

接入上游模型

建议将模型调用封装为独立生产者:调用开始时写入 start 事件,收到增量内容时写入 delta 事件,结束时写入 doneerror 事件。业务层只消费统一事件,不直接解析某家服务的响应格式。

伪代码结构如下:

async def run_model(task_id: str, prompt: str) -> None:
    await append_event(task_id, EventIn(kind="start", data={
   }))
    try:
        async for chunk in upstream_stream(prompt):
            await append_event(task_id, EventIn(
                kind="delta", data={
   "text": chunk}
            ))
        await append_event(task_id, EventIn(kind="done", data={
   }))
    except Exception as exc:
        await append_event(task_id, EventIn(
            kind="error", data={
   "code": "upstream_failed"}
        ))
        raise

真实实现还应加入请求超时、取消传播、重试边界和响应校验。对于不可幂等的上游请求,不要因为网络超时就无条件重试,否则可能产生重复计费或重复执行。日志中记录任务 ID、模型标识、耗时和错误类别即可,提示词和输出内容是否落盘应由数据分级和合规要求决定。

反向代理配置

如果前面使用 Nginx,需要关闭响应缓冲,并设置合理的读取超时。配置只是示例,具体值要结合最长任务时间和网络环境调整:

location /tasks/ {
   
    proxy_pass http://127.0.0.1:8000;
    proxy_http_version 1.1;
    proxy_set_header Connection "";
    proxy_buffering off;
    proxy_cache off;
    proxy_read_timeout 1h;
}

还要检查负载均衡器是否会主动截断长连接,以及入口层是否限制响应时间。心跳只能帮助发现连接状态,不能替代事件持久化。

常见问题

断线后会不会重复显示?

会有这种可能。网络发送成功与客户端业务处理成功不是同一个时刻。推荐让客户端按事件 ID 去重,并把最后一个已处理 ID 作为重连游标。若只要求最终结果,客户端也可以在 done 事件到达后重新拉取完整结果。

为什么不直接用 Redis Pub/Sub?

Pub/Sub 更适合在线广播,订阅者断开期间通常无法从频道补回历史消息。需要断线续传时,应选用带持久化记录和游标的机制,或自行建立事件表。

多个订阅者是否应该共享一个模型请求?

通常应该共享上游任务,而不是让每个订阅者各自调用模型。这样可以避免重复计算,但要明确任务权限:订阅接口必须校验用户是否有权读取对应任务。事件存储也应设置保留期限,避免敏感内容长期存在。

SSE 能否双向通信?

SSE 是服务端到客户端的单向通道。客户端提交提示词、取消任务或发送确认,应使用普通 HTTP、WebSocket 或其他明确的双向协议。不要把控制消息偷偷塞进 SSE 数据流。

如何限制慢客户端?

生产实现需要跟踪每个连接的发送队列长度、最近写入时间和异常次数。超过阈值时断开连接,让客户端依据游标重连;同时保留任务事件,避免通过内存队列无限等待。若事件量很大,应把实时订阅和离线消费拆成不同通道。

总结

流式模型接入的关键不在于把响应拆成更小的字符串,而在于建立可恢复的事件边界。事件先持久化、订阅者按游标读取、客户端按 ID 去重,再配合连接超时、慢消费者隔离和审计字段,才能把一次性演示提升为可维护的服务。

Socket 和文件描述符解释了连接如何被操作系统管理,SSE 提供了浏览器友好的传输方式,而事件日志承担了恢复与追踪职责。三者职责不同,分层后才能替换上游模型、扩展订阅端,并在出现异常时定位究竟是模型、存储、代理还是客户端环节发生了问题。

相关文章
|
安全 搜索推荐 网络架构
什么是内网和外网?什么是内网IP和外网IP?本地连接和宽带连接又有什么区别?
何为内网外网迷糊?究竟什么是内网?什么是外网?他们又有和区别?还有什么是内网IP和外网IP?本地连接和宽带连接有什么区别?怂怂今天就来给大家科普一下吧:
16310 0
|
Java 数据库连接
什么是双亲委派?如何打破双亲委派?
什么是双亲委派?如何打破双亲委派?
599 0
|
Java Spring 容器
spring之HttpInvoker
  一、HttpInvoker是常用的Java同构系统之间方法调用实现方案,是众多Spring项目中的一个子项目。顾名思义,它通过HTTP通信即可实现两个Java系统之间的远程方法调用,使得系统之间的通信如同调用本地方法一般。
2925 0
|
应用服务中间件 索引 nginx
生产环境ES查询延迟排查
最近经常收到业务方配置的ES查询延迟告警,同样的请求手动在Kibana控制台执行只需几十毫秒就返回结果。受影响的整个链路情况如下,php应用程序通过部署在ES集群各节点上的nginx访问ES请求查询数据。
6073 0
|
3月前
|
人工智能 安全 JavaScript
阿里云无影AgentBay对接全指南:MCP/SDK/Web全链路接入与实战
在AI智能体快速落地的当下,安全、稳定、可扩展的云端执行环境成为核心刚需。阿里云无影AgentBay作为专为AI Agent打造的云端沙箱基础设施,提供浏览器、桌面、代码、移动端四大场景的隔离执行能力,解决了本地环境依赖、安全风险、并发限制等痛点,是构建企业级智能体的首选底座。2026年,AgentBay已完成多轮迭代,接入方式更灵活、环境更丰富、生态更完善,支持MCP协议、多语言SDK、Web SDK三种主流接入方式,覆盖从简单工具调用到复杂自动化流程的全场景需求。本文从核心概念、接入准备、三种对接方式、实战案例、高级配置到运维优化,提供完整的对接使用指南,搭配可直接运行的代码命令,帮助开发
559 4
|
5月前
|
机器学习/深度学习 人工智能 Java
实测对比:企业落地的主流 AI 开发框架测评
本文以中立、客观、可落地为原则,实测对比JBoltAI、LangChain、Spring AI等主流AI框架,聚焦Java企业适配性、国产模型支持、工程化能力及存量系统改造难度,提供清晰选型建议。(239字)
365 0
|
设计模式 消息中间件 监控
并发设计模式实战系列(5):生产者/消费者
🌟 ​大家好,我是摘星!​ 🌟今天为大家带来的是并发设计模式实战系列,第五章,废话不多说直接开始~
454 1
|
Java 开发者 Spring
java springboot监听事件和处理事件
通过上述步骤,开发者可以在Spring Boot项目中轻松实现事件的发布和监听。事件机制不仅解耦了业务逻辑,还提高了系统的可维护性和扩展性。掌握这一技术,可以显著提升开发效率和代码质量。
627 13
|
负载均衡 Dubbo Java
哈啰面试:说说Dubbo运行原理?
哈啰面试:说说Dubbo运行原理?
338 0
哈啰面试:说说Dubbo运行原理?
|
Java C语言
STM32使用printf重定向到USART(串口)并打印数据到串口助手
STM32使用printf重定向到USART(串口)并打印数据到串口助手
2544 0