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

event-source-polyfill无法获取Rust SSE API数据的排查求助

问题描述

我用Rust的Actix-Web框架实现了一个SSE API,用curl -N http://localhost:8000/sse/producer测试能正常返回数据,但在React项目中用event-source-polyfill(版本^1.0.31)编写的客户端无法获取数据,日志仅显示无有效信息的错误。后续尝试改用标准格式的SSEMessage结构体发送消息,问题依然存在。

服务端初始代码、Cargo依赖、客户端代码及代理配置如下:

服务端初始代码

use actix_web::{App, HttpResponse, HttpServer, Responder, post, get};
use rust_wheel::common::util::net::sse_stream::SseStream;
use std::time::{Duration, SystemTime};
use tokio::{
    sync::mpsc::{UnboundedReceiver, UnboundedSender},
    task,
};

#[get("/sse/producer")]
pub async fn sse() -> impl Responder {
    let (tx, rx): (UnboundedSender<String>, UnboundedReceiver<String>) =
        tokio::sync::mpsc::unbounded_channel();
    task::spawn(async move {
        for _ in 0..5 {
            let message = format!("Current time: {:?}", SystemTime::now());
            tx.send(message).unwrap();
            tokio::time::sleep(Duration::from_secs(1)).await;
        }
    });
    let response = HttpResponse::Ok()
        .content_type("text/event-stream")
        .streaming(SseStream { receiver: Some(rx) });
    response
}

#[actix_web::main]
async fn main() -> Result<(), std::io::Error> {
    HttpServer::new(|| App::new().service(sse).service(sse))
        .bind("127.0.0.1:8000")?
        .run()
        .await
}

Cargo.toml依赖

tokio = { version = "1.17.0", features = ["full"] }
serde = { version = "1.0.64", features = ["derive"] }
serde_json = "1.0.64"
rust_wheel = { git = "https://github.com/jiangxiaoqiang/rust_wheel.git", branch = "diesel2.0" }
actix-web = "4"

React客户端代码

import React from 'react';
import { EventSourcePolyfill } from 'event-source-polyfill';

const App: React.FC = () => {

  React.useEffect(() => {
    doSse();
  },[]);

  const doSse = () => {
    let eventSource: EventSourcePolyfill;    
    eventSource = new EventSourcePolyfill('/sse/producer', {
      headers: {
      }
    });
    eventSource.onopen = () => {
      
    }
    eventSource.onerror = (error:any) => {
      console.log(error)
      eventSource.close();
    }
    eventSource.onmessage = (msg: any) => {
      console.log(msg)
    };
  }

  return (
    <div>
    </div>
  );
}

export default App;

客户端代理配置

proxy: {
      '/sse': 'http://localhost:8000/'
    }
修复方案

1. 严格遵循SSE格式规范

SSE要求每条消息必须以data: 前缀开头,且以两个换行符(\n\n)结尾。你的SseStream大概率没有正确格式化输出:

  • 检查rust_wheel中的SseStream实现,确保它将SSEMessage序列化为标准格式,可参考如下逻辑:
    use std::fmt;
    
    #[derive(Deserialize, Serialize)]
    pub struct SSEMessage {
        pub event: Option<String>,
        pub data: String,
        pub id: Option<String>,
        pub retry: Option<u32>,
    }
    
    impl fmt::Display for SSEMessage {
        fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
            // 输出event字段(可选)
            if let Some(event) = &self.event {
                writeln!(f, "event: {}", event)?;
            }
            // 多行data需每行添加data: 前缀
            for line in self.data.split('\n') {
                writeln!(f, "data: {}", line)?;
            }
            // 输出id字段(可选)
            if let Some(id) = &self.id {
                writeln!(f, "id: {}", id)?;
            }
            // 输出retry字段(可选)
            if let Some(retry) = self.retry {
                writeln!(f, "retry: {}", retry)?;
            }
            // 必须以空行结束消息
            writeln!(f)
        }
    }
    
  • 如果不想修改rust_wheel,初始版本中可直接拼接符合规范的字符串:
    let message = format!("data: Current time: {:?}\n\n", SystemTime::now());
    

2. 完善Actix-Web响应头

添加防止缓存、保持连接的响应头,解决跨域或缓存问题:

let response = HttpResponse::Ok()
    .content_type("text/event-stream")
    .insert_header(("Cache-Control", "no-cache"))
    .insert_header(("Connection", "keep-alive"))
    .insert_header(("Access-Control-Allow-Origin", "*")) // 开发环境临时使用,生产环境请指定具体域名
    .streaming(SseStream { receiver: Some(rx) });

3. 调整React客户端配置与错误处理

  • 优化错误处理,不要直接关闭连接,先打印详细错误信息:
    eventSource.onerror = (error) => {
      console.error('SSE错误详情:', error);
      console.error('当前连接状态:', eventSource.readyState);
      // 仅在连接状态为CLOSED时关闭,避免误杀
      if (eventSource.readyState === EventSource.CLOSED) {
        eventSource.close();
      }
    };
    
  • 更新代理配置,支持长连接传递:
    proxy: {
      '/sse': {
        target: 'http://localhost:8000/',
        changeOrigin: true,
        ws: true,
        onProxyReq: (proxyReq) => {
          proxyReq.setHeader('Connection', 'keep-alive');
        }
      }
    }
    

4. 修复Actix-Web服务重复注册问题

main函数中重复注册了sse服务,会导致路由冲突,改为:

#[actix_web::main]
async fn main() -> Result<(), std::io::Error> {
    HttpServer::new(|| App::new().service(sse))
        .bind("127.0.0.1:8000")?
        .run()
        .await
}

5. 增加消息发送错误处理

替换unwrap()为错误日志,避免通道异常导致消息丢失:

if let Err(e) = tx.send(message) {
    eprintln!("发送SSE消息失败: {}", e);
}

内容的提问来源于stack exchange,提问作者Dolphin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 08:27:11