如何在Actix Web中实现Server-Sent Events流式传输OpenAI响应
用Actix Web实现OpenAI流式响应的Server-Sent Events传输
要把OpenAI的流式响应通过SSE推送给客户端,你需要调整handler逻辑,用Actix Web Lab的SSE模块构建流式响应,替代原来的普通HTTP响应。以下是修改后的完整代码:
use actix_web::{get, web, App, HttpResponse, HttpServer, Responder}; use actix_web_lab::sse; use async_stream::stream; use futures::stream::StreamExt; use reqwest; use reqwest_eventsource::{Event, EventSource}; use serde_json::Value; use std::convert::Infallible; // 替换成你的OpenAI API地址和密钥 const OPEN_AI_COMPLETION_URL: &str = "https://api.openai.com/v1/chat/completions"; const OPEN_AI_KEY: &str = "your-api-key-here"; #[get("/")] async fn stream_openai_response() -> sse::Response { // 构建OpenAI流式请求 let request = reqwest::Client::new() .post(OPEN_AI_COMPLETION_URL) .header("Authorization", format!("Bearer {}", OPEN_AI_KEY)) .json(&serde_json::json!({ "model": "gpt-3.5-turbo", "messages": [ {"role": "system", "content": "History of USA?"}, ], "max_tokens": 100, "temperature": 0, "stream": true, })); // 创建EventSource处理OpenAI的流式响应 let mut es = match EventSource::new(request) { Ok(es) => es, Err(err) => { return sse::Response::new(stream! { yield sse::Event::Data(err.to_string()); }) } }; // 创建SSE流,将OpenAI的消息转发给客户端 let sse_stream = stream! { while let Some(event) = es.next().await { match event { Ok(Event::Message(message)) => { // 跳过空消息 if message.data.is_empty() { continue; } // 处理OpenAI的结束信号 if message.data == "[DONE]" { break; } // 解析JSON获取content字段 if let Ok(Value::Object(obj)) = serde_json::from_str(&message.data) { if let Some(choices) = obj.get("choices") { if let Some(choice) = choices.as_array().and_then(|arr| arr.first()) { if let Some(delta) = choice.get("delta") { if let Some(content) = delta.get("content") { if let Some(content_str) = content.as_str() { // 发送内容作为SSE数据事件 yield sse::Event::Data(content_str.to_string()); } } } } } } } Ok(Event::Open) => { // 可选:连接建立时发送通知 yield sse::Event::Data("Connection to OpenAI established!".to_string()); } Err(err) => { // 发送错误信息并结束流 yield sse::Event::Data(format!("Error: {}", err)); break; } } } }; // 返回SSE响应 sse::Response::new(sse_stream) } #[actix_web::main] async fn main() -> std::io::Result<()> { HttpServer::new(|| { App::new() .service(stream_openai_response) }) .bind(("127.0.0.1", 8080))? .run() .await }
关键说明:
- SSE响应构建:用
sse::Response::new()包装异步流,Actix Web会自动处理SSE的响应头(Content-Type: text/event-stream等)。 - OpenAI流解析:OpenAI的流式响应每条消息是JSON格式,需要解析出
choices[0].delta.content字段,这才是实际的文本内容。 - 结束信号处理:当收到
data: [DONE]时,跳出循环结束流,客户端会收到流结束的通知。 - 错误处理:创建EventSource失败或中间出错时,会向客户端发送错误信息并关闭流。
客户端测试:
你可以用浏览器访问http://localhost:8080,打开开发者工具的控制台查看SSE消息;或者用curl命令:
curl -N http://localhost:8080
内容的提问来源于stack exchange,提问作者Nane
相关产品推荐
相关产品推荐

