功能和问题
- 功能描述: 在ai开发中,需要流式返回的场景,可以使用StreamingResponse来实现。 他的用法简单,只要返回一个可迭代对象,就可以实现流式响应。
代码片段展示
在router中,使用StreamingResponse来实现流式响应。
@router.post("/chat")
def chat():
return StreamingResponse(
stream_message(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no",
"Connection": "keep-alive",
},
)
重点是要返回一个可迭代对象,例如generator函数,当我们在大模型交互中,可以使用类似的stream的api调用大模型的接口
def stream_message():
for n in range(10):
time.sleep(1)
yield f"message {n}"
前端接收方式
方式一:EventSource
原生 SSE 接收方式,自动重连、使用简单,但只支持 GET 且无法自定义请求头。
const sse = new EventSource("/chat")
sse.onopen = () => {
console.log("connected")
}
sse.onmessage = (event) => {
console.log(event.data)
}
sse.onerror = () => {
sse.close()
}
方式二:Fetch API
适合需要 POST、自定义 header 的场景,需要自己解析 SSE 的 data 行。
const response = await fetch("/chat", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ prompt: "hi" })
})
const reader = response.body.getReader()
const decoder = new TextDecoder("utf-8")
let buffer = ""
while (true) {
const { value, done } = await reader.read()
if (done) break
buffer += decoder.decode(value, { stream: true })
const parts = buffer.split("\n\n")
buffer = parts.pop() || ""
for (const part of parts) {
const line = part.split("\n").find((item) => item.startsWith("data:"))
if (line) {
const data = line.replace(/^data:\s*/, "")
console.log(data)
}
}
}
方式三:fetch-event-source
第三方库封装 SSE,支持 POST、自定义 header、自动重连,使用前需安装依赖。
import { fetchEventSource } from "@microsoft/fetch-event-source"
await fetchEventSource("/chat", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ prompt: "hi" }),
onmessage(event) {
console.log(event.data)
},
onerror(err) {
throw err
}
})