You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.10 12:18:13