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
相关产品推荐
相关产品推荐

