SSE 流式传输实现方法详解(AI 视频下载总结器的实时响应)

体验瓶颈的第一名是”等待”,而不是模型选型。用户粘贴一个视频链接,从下载、字幕提取到 AI 总结,整条链路要跑 30 秒到一分钟,如果页面只能转圈干等,很多人中途就关掉走了。我先后对比过轮询、WebSocket 和 SSE,最终选了 SSE(Server-Sent Events):它跑在 HTTP 之上,服务端可以边算边推,实现成本比 WebSocket 低一大截,前端用原生 EventSource 或者 fetch 的 ReadableStream 就能消费。项目上线后首字延迟稳定在 500ms 以内,用户从提交链接到看到第一行总结基本没有等待感。下面把我在项目里把 SSE 落地的完整过程拆开讲,包括协议细节、代码实现和线上踩过的坑。

一、为什么用 SSE

AI 总结是一个”长任务”:下载视频(5-30s)→ 字幕提取(5-10s)→ AI 总结(10-30s)。如果同步等结果返回,用户体验差。

一开始我图省事,先写了个同步接口,前端轮询任务状态。结果发现两个问题:轮询有固定间隔,进度永远慢半拍;任务没跑完时连接一直挂着,白白占资源。SSE 的思路完全不同,服务端把状态变化主动推给客户端,客户端不用反复问。我在项目里按事件把整个流程拆开:

  • 客户端立即收到”开始处理”事件;
  • 阶段进度实时推送(下载中/转录中/总结中);
  • AI 总结的 token 逐字推送(流式);
  • 完成时收到”完成”事件。

这套事件设计后来帮我省了不少事:stage 事件驱动进度条,token 事件负责逐字打字效果,done 事件收尾。前端只需要按事件名各干各的,逻辑非常干净。

二、SSE 协议

决定用 SSE 之后,我先把协议规范通读了一遍。SSE 的报文就是几行纯文本,用空行把一条条消息隔开,简单到可以直接用眼睛读:

event: message
data: {"text": "第一段"}

event: message
data: {"text": "第二段"}

event: done
data: {"ok": true}

每条消息:

  • event:事件类型(自定义);
  • data:数据(字符串);
  • id:消息 ID(可选,用于断点续传);
  • retry:客户端重连间隔(毫秒)。

协议简单,但坑都在细节里。比如每条消息必须用双换行收尾,客户端才能判定一条完整事件;数据里的空行要拆成多行 data 字段。这些我调试时都踩过,后面会单独说。

响应头是让 SSE 真正”流”起来的关键:

Content-Type: text/event-stream
Cache-Control: no-cache
Connection: keep-alive
X-Accel-Buffering: no  # 关键:禁用 Nginx 缓冲

这四个头缺一不可,尤其 X-Accel-Buffering: no。不设置它,Nginx 会把流攒成一大坨再吐给前端,实时性全没了,这个坑在部署章节还会展开。

三、FastAPI 集成

服务端我用 FastAPI 承接,社区有现成的 sse-starlette,几行代码就能把异步生成器包装成 SSE 响应,不用自己处理协议层。

pip install sse-starlette
from sse_starlette.sse import EventSourceResponse
import asyncio

@app.post('/api/summary/stream')
async def stream_summary(req: SummaryRequest):
    """流式返回视频总结"""
    async def event_generator():
        # 1) 阶段 1:开始下载
        yield {
            'event': 'stage',
            'data': json.dumps({'stage': 'downloading', 
                                'percent': 0})
        }
        
        # 2) 下载视频
        info = await downloader.download(req.url, 
                                          progress_callback=lambda p: 
                                              yield_progress(p))
        yield {
            'event': 'stage',
            'data': json.dumps({'stage': 'transcribing', 
                                'percent': 0})
        }
        
        # 3) 提取字幕
        subtitle = await transcriber.extract(info)
        yield {
            'event': 'stage',
            'data': json.dumps({'stage': 'summarizing', 
                                'percent': 0})
        }
        
        # 4) AI 总结(流式)
        async for chunk in ai_service.stream_summarize(subtitle):
            yield {
                'event': 'token',
                'data': json.dumps({'text': chunk})
            }
        
        # 5) 完成
        yield {
            'event': 'done',
            'data': json.dumps({'summary_id': 'xxx'})
        }
    
    return EventSourceResponse(event_generator())

这段代码的核心是 event_generator 这个异步生成器,它按阶段依次 yield 事件,一次覆盖下载、转录、总结的完整生命周期。要注意的是生成器里每个 yield 都必须尽快返回,任何阻塞调用都会卡住整条流,所以下载和 AI 调用一律走异步客户端。

四、DeepSeek 流式调用

字幕拿到手,重头戏是 AI 流式总结。DeepSeek 的接口兼容 OpenAI 协议,把 stream 打开后,模型会把结果切成一块块增量推回来:

from openai import AsyncOpenAI

client = AsyncOpenAI(
    api_key=settings.DEEPSEEK_API_KEY,
    base_url='https://api.deepseek.com'
)

async def stream_summarize(self, text: str) -> AsyncIterator[str]:
    stream = await client.chat.completions.create(
        model='deepseek-chat',
        messages=[
            {'role': 'system', 'content': '你是视频总结助手...'},
            {'role': 'user', 'content': f'总结:\n{text[:4000]}'}
        ],
        stream=True
    )
    
    async for chunk in stream:
        delta = chunk.choices[0].delta.content
        if delta:
            yield delta

我把字幕截到 4000 字符再进模型,超过的部分走分段逻辑,避免一次请求塞太多输入拖慢首字时间。实测开流之后,第一块 token 大约 300ms 就能回到服务端,这是前端”逐字出字”效果的来源。

五、进度推送

阶段进度有了,但 AI 总结那十几秒里,用户盯着一个”总结中”还是空落落的。我用了 token 计数反推百分比的笨办法,零额外开销:

async def stream_summarize(self, text: str, total_tokens: int):
    stream = await client.chat.completions.create(
        model='deepseek-chat',
        messages=[...],
        stream=True
    )
    
    produced = 0
    async for chunk in stream:
        delta = chunk.choices[0].delta.content
        if delta:
            produced += len(delta)  # 估算
            yield delta
            
            # 推送进度
            yield_progress({
                'percent': min(100, produced * 100 // total_tokens),
                'produced': produced
            })

这个方法精度不高,胜在简单——不用额外接口,也不让模型多说一句废话。前端配合平滑动画,百分比从 0 爬到 100,观感是顺的。

六、前端处理 SSE

1. 原生 EventSource

前端最简单的做法是直接用 EventSource,它内置自动重连,按事件名监听就能收到推送:

const es = new EventSource('/api/summary/stream', {
    withCredentials: true  // 携带 cookie
})

es.addEventListener('stage', (e) => {
    const { stage, percent } = JSON.parse(e.data)
    updateProgress(stage, percent)
})

es.addEventListener('token', (e) => {
    const { text } = JSON.parse(e.data)
    summaryDiv.textContent += text  // 追加显示
})

es.addEventListener('done', (e) => {
    const { summary_id } = JSON.parse(e.data)
    es.close()
    showComplete(summary_id)
})

es.addEventListener('error', (e) => {
    if (es.readyState === EventSource.CLOSED) return
    console.error('SSE 错误', e)
})

用下来我发现 EventSource 有硬伤:

  • 只支持 GET;
  • 不能自定义请求体(POST 不行);
  • 不能设自定义 Header。

我们的接口要 POST 视频链接和参数,还要带鉴权头,原生 EventSource 直接出局。

2. fetch + ReadableStream(推荐)

AI 视频下载要传 URL+参数,POST 请求,EventSource 不行。平台用 fetch + ReadableStream:

async function streamSummary(url, body) {
    const response = await fetch('/api/summary/stream', {
        method: 'POST',
        headers: { 'Content-Type': 'application/json' },
        body: JSON.stringify(body)
    })
    
    const reader = response.body.getReader()
    const decoder = new TextDecoder()
    let buffer = ''
    
    while (true) {
        const { done, value } = await reader.read()
        if (done) break
        
        buffer += decoder.decode(value, { stream: true })
        
        // SSE 按 \n\n 分隔事件
        const events = buffer.split('\n\n')
        buffer = events.pop()  // 最后一段可能不完整
        
        for (const evt of events) {
            const lines = evt.split('\n')
            let eventType = 'message'
            let data = ''
            for (const line of lines) {
                if (line.startsWith('event:')) eventType = line.slice(6).trim()
                else if (line.startsWith('data:')) data += line.slice(5).trim()
            }
            handleEvent(eventType, JSON.parse(data))
        }
    }
}

这段我调了一晚上才稳。核心是 buffer 的拼接逻辑:每次读到的字节先拼进 buffer,再按双换行符切事件,最后一段拼不完整的留在 buffer 里等下一批数据。半包、拆包这些边界情况全测一遍之后,再也没出现过丢事件或者 JSON 解析失败。

七、AbortController 取消

流式传输还有一个绕不开的问题:用户不想等了。浏览器端用 AbortController 随时可以掐断 fetch:

let controller = null

function startSummary() {
    controller = new AbortController()
    fetch('/api/summary/stream', {
        method: 'POST',
        signal: controller.signal,
        // ...
    })
}

function cancelSummary() {
    controller.abort()
}

前端掐断只是第一步,后端也得感知连接断开,否则模型还在继续算,白白烧 token。后端捕获取消:

from starlette.requests import ClientDisconnect

async def event_generator():
    try:
        async for chunk in ai.stream(...):
            yield {...}
    except asyncio.CancelledError:
        log.info('用户取消')
        raise

后端依赖 asyncio.CancelledError 做清理,配合日志把”用户取消”和”系统异常”区分开,排查问题时一眼就知道是哪一类。

八、断线重连

网络抖动时 SSE 会掉线。原生 EventSource 自带自动重连,但 fetch 方案没有,得自己补:

async function streamWithRetry(url, body, maxRetry = 3) {
    let attempt = 0
    while (attempt < maxRetry) {
        try {
            await streamSummary(url, body)
            return  // 成功
        } catch (e) {
            if (e.name === 'AbortError') return  // 用户取消不重试
            attempt++
            await sleep(1000 * attempt)  // 退避
        }
    }
    showError('网络异常,请重试')
}

重试我做了两层保护:用户主动取消一律不重试,避免误弹提示;指数退避从 1 秒开始递增,防止断网恢复瞬间所有客户端同时重连,把服务端打挂。

九、Nginx 关键配置

代码全对了,部署上线时又栽了一个跟头——Nginx 默认会给代理响应开缓冲。SSE 必须禁用:

location /api/stream/ {
    proxy_pass http://backend;
    proxy_http_version 1.1;
    proxy_buffering off;       # 关键
    proxy_cache off;
    proxy_set_header Connection '';
    chunked_transfer_encoding on;
    proxy_read_timeout 600s;   # 长连接超时
}

当时线上首字延迟飙到 30 秒以上,排查半天才发现是 proxy_buffering 没关。关掉缓冲、把超时拉到 600 秒之后,实时性立刻恢复。不设 proxy_buffering off 会有 30s+ 延迟。这个坑我建议直接写进部署文档,换环境时少踩一次。上线前按这三步核对:

  1. 确认 proxy_buffering off 和 proxy_cache off 都已生效;
  2. 确认 proxy_read_timeout 大于最长任务耗时;
  3. 用 curl 观察输出,确认数据是一块块到达而不是整块吐出。

十、错误处理

流式接口的错误处理和普通接口不一样:流可能已经开推了,这时再返回 HTTP 状态码没有意义,所以错误必须作为事件推进去:

async def event_generator():
    try:
        yield {'event': 'stage', 'data': json.dumps(
            {'stage': 'downloading'})}
        info = await downloader.download(req.url)
    except DownloadError as e:
        yield {'event': 'error', 'data': json.dumps({
            'code': 'DOWNLOAD_FAILED',
            'message': str(e)
        })}
        return
    except Exception as e:
        log.exception('未知错误')
        yield {'event': 'error', 'data': json.dumps({
            'code': 'INTERNAL_ERROR',
            'message': '系统异常'
        })}
        return
    
    try:
        # ...继续
    except Exception:
        # 错误也要 yield
        ...

前端统一监听 error 事件,把 code 映射成中文提示。下载失败和 AI 异常分开给码,用户能清楚知道问题出在哪一步,我们也能按 code 聚合告警。

十一、超时控制

SSE 是长连接,不设最大时长就是隐患。万一某个视频解析卡死,连接会一直挂着占资源:

async def event_generator():
    start = time.time()
    max_duration = 300  # 5 分钟
    
    async for chunk in ai.stream(...):
        if time.time() - start > max_duration:
            yield {'event': 'error', 'data': json.dumps({
                'code': 'TIMEOUT',
                'message': '处理超时'
            })}
            return
        yield chunk

5 分钟这个上限我压过几轮,正常视频 1-3 分钟跑完,超时多发生在资源异常的链接上,直接掐掉合理。

十二、与 WebSocket 对比

项目里有人建议直接用 WebSocket,我当时把两个方案摊开对比了一遍:

维度 SSE WebSocket
方向 服务端→客户端 双向
协议 HTTP WS
重连 自动 手动
复杂度 低 中
适用 单向流 双向通信

结论很清楚:我们的场景是服务端向客户端单向推送,WebSocket 的双向能力用不上,还要多维护一套协议和心跳,复杂度不划算。AI 视频下载是”服务端→客户端”单向,SSE 完美匹配。

十三、性能监控

上线后我加了连接级埋点,至少要知道同时有多少条流在跑、平均耗时多少:

import time

metrics = {
    'sse_connections_active': 0,
    'sse_total_duration': 0
}

async def event_generator():
    start = time.time()
    metrics['sse_connections_active'] += 1
    try:
        # ... 业务逻辑
    finally:
        metrics['sse_total_duration'] = time.time() - start
        metrics['sse_connections_active'] -= 1

这些指标接进监控面板,配合前面的事件日志,高峰期该扩容还是该优化代码,一看就有数。

十四、踩过的坑

这一节把线上真实踩过的坑汇总一下,都是血泪教训:

  • Nginx 缓冲:不关 proxy_buffering off 延迟 30s+。
  • HTTP/1.1 6 并发限制:浏览器同域最多 6 个并发 HTTP/1.1。AI 总结要排队。
  • EventSource GET 限制:传 URL 用 GET,无 body;POST 用 fetch。
  • 断线重连风暴:用户断网后 EventSource 自动重连可能打挂服务。限制重连间隔。
  • CORS preflight:POST + custom header 触发 OPTIONS 预检。预检通过后才能 SSE。
  • Heartbeat:长连接可能被中间代理掐断。定期发心跳(如 : heartbeat\n\n)。
  • 压缩:gzip 压缩 SSE 流要小心,可能让 event: 边界错乱。

七个坑里,Nginx 缓冲和断线重连风暴影响面广,建议代码评审时重点查这两条;心跳和压缩是部署在特殊网络环境时才需要关注。

十五、与 WebSocket 取舍

回到最初的选型问题,我给团队的理由很直白:

  • AI 总结单向推送;
  • SSE 协议简单(HTTP 基础上);
  • 自动重连(EventSource);
  • Nginx 配置简单。

如果未来要做实时协作编辑、多端同步这类双向通信,再引入 WebSocket 不迟。现阶段 SSE 已经把我们单向前推送的场景吃满了,架构上也给切换留了余地。

常见问题(FAQ)

Q1:SSE 能传文件吗?

能,但效率不如 WebSocket。AI 总结是文本流,SSE 合适。

Q2:SSE 跨域怎么配?

同 CORS,但要允许 text/event-stream Content-Type。

Q3:SSE 浏览器兼容性?

IE 不支持,Edge/Chrome/Firefox/Safari 都支持。平台目标用户不用 IE。

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 qiqicto@qq.com 举报,一经查实,本站将立刻删除。
赞 (0)
其AI的头像其AI普通用户

相关推荐

返回顶部