Java SSE案例

wen java案例 1

Java SSE 服务端推送实战:从入门到高并发优化(附完整案例)

目录导读

  • 什么是SSE?与WebSocket有什么区别?
  • Java实现SSE的三种核心方式(Spring MVC / JAX-RS / 原生Servlet)
  • 实战案例:基于Spring Boot构建实时股票行情推送
  • SSE连接保活、断线重连与心跳机制
  • 高并发场景下的SSE性能调优与注意事项
  • 常见问题问答(FAQ)

什么是SSE?与WebSocket有什么区别?

SSE(Server-Sent Events,服务端发送事件)是一种基于HTTP的轻量级实时通信协议,它允许服务器通过单个长连接,主动向浏览器持续推送数据,而客户端只需使用原生EventSource对象即可接收。

Java SSE案例

与WebSocket的核心区别

  • 单向 vs 双向:SSE是服务器→客户端单向推送;WebSocket是双向全双工通信。
  • 协议复杂度:SSE基于普通HTTP,无需额外握手协议;WebSocket需要升级协议。
  • 自动重连:SSE内置断线重连机制(浏览器自动处理);WebSocket需手动实现。
  • 适用场景:SSE适合新闻推送、股票行情、日志流、AI对话流式输出;WebSocket适合聊天、在线游戏、协同编辑。

一句话总结:如果只需要服务端推数据,SSE比WebSocket更简单、更可靠、更省资源。


Java实现SSE的三种核心方式

Spring MVC + SseEmitter(最常用)

Spring 4.2+原生支持SSE,通过返回SseEmitter对象即可。

@RestController
public class StockController {
    private final CopyOnWriteArrayList<SseEmitter> emitters = new CopyOnWriteArrayList<>();
    @GetMapping("/stock/stream")
    public SseEmitter stream() {
        SseEmitter emitter = new SseEmitter(0L); // 0表示超时无限
        emitters.add(emitter);
        emitter.onCompletion(() -> emitters.remove(emitter));
        emitter.onTimeout(() -> emitters.remove(emitter));
        return emitter;
    }
    // 推送线程:模拟行情数据
    @Scheduled(fixedRate = 1000)
    public void pushStockData() {
        String data = "{\"code\":\"AAPL\",\"price\":" + Math.random() * 100 + "}";
        for (SseEmitter emitter : emitters) {
            try {
                emitter.send(SseEmitter.event().name("stock").data(data));
            } catch (IOException e) {
                emitters.remove(emitter);
            }
        }
    }
}

JAX-RS(Jersey) + SseEventSink

适用于标准Java EE环境,使用SseEventSinkSse上下文。

原生Servlet + AsyncContext

通过AsyncContext维持连接,手动写入response.getOutputStream(),需自行处理编码和心跳。


实战案例:基于Spring Boot构建实时股票行情推送

需求:前端页面实时显示5只股票的价格,每秒更新一次。

后端完整实现(含线程池管理)

@Service
public class StockService {
    private final Map<String, Double> stockPrices = new ConcurrentHashMap<>();
    private final List<SseEmitter> clients = new CopyOnWriteArrayList<>();
    public SseEmitter subscribe() {
        SseEmitter emitter = new SseEmitter(30_000L); // 30秒超时
        clients.add(emitter);
        emitter.onTimeout(() -> clients.remove(emitter));
        emitter.onError((e) -> clients.remove(emitter));
        return emitter;
    }
    @Scheduled(fixedRate = 1000)
    public void updatePrices() {
        stockPrices.forEach((code, price) -> {
            double newPrice = price + (Math.random() - 0.5) * 2;
            stockPrices.put(code, newPrice);
            broadcast(code, newPrice);
        });
    }
    private void broadcast(String code, double price) {
        String payload = String.format("{\"code\":\"%s\",\"price\":%.2f}", code, price);
        for (SseEmitter emitter : clients) {
            try {
                emitter.send(SseEmitter.event().name("quote").data(payload));
            } catch (IOException e) {
                clients.remove(emitter);
            }
        }
    }
}

前端HTML(原生EventSource)

<!DOCTYPE html>
<html>
<body>
<div id="quotes"></div>
<script>
    const source = new EventSource('/stock/stream');
    source.addEventListener('quote', (event) => {
        const data = JSON.parse(event.data);
        document.getElementById('quotes').innerHTML += 
            `<p>${data.code}: $${data.price.toFixed(2)}</p>`;
    });
    source.onerror = () => console.log('连接断开,浏览器自动重连...');
</script>
</body>
</html>

关键点SseEmitter.event().name("quote")指定了事件类型,前端用addEventListener监听;若使用data()而不指定事件名,前端用onmessage接收。


SSE连接保活、断线重连与心跳机制

心跳保活(防止代理超时断开)

默认SSE空闲超过60秒可能被Nginx等代理断开,解决方案:

  • 后端每15秒发送一个注释行()作为心跳。
    emitter.send(SseEmitter.event().comment("heartbeat"));
  • 或者发送无数据事件:emitter.send(SseEmitter.event().data(""));

断线重连(浏览器内置)

EventSource自动重连,默认间隔3秒,可通过HTTP响应头Retry-After自定义重连间隔。

emitter.send(SseEmitter.event().reconnectTime(5000));

最佳实践

  • 设置合理的超时时间,建议30秒~2分钟,配合心跳保活。
  • SseEmitter完成、超时、异常回调中必须清理资源,防止内存泄漏。
  • 使用CopyOnWriteArrayList存储emitter,保证并发安全且不阻塞读。

高并发场景下的SSE性能调优

问题 解决方案
大量闲置连接占用线程 启用Tomcat/Nginx的async supported,配置合适的maxConnections
服务端广播压力大 将推送任务放入独立线程池,避免阻塞业务线程
前端反馈数据丢失 开启压缩(gzip),并检查SSE格式是否包含data:和双换行
连接数超出限制 结合Redis发布订阅,实现多实例SSE水平扩展

示例:使用ThreadPoolTaskExecutor处理推送任务

@Bean
public ThreadPoolTaskExecutor sseExecutor() {
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setCorePoolSize(4);
    executor.setMaxPoolSize(16);
    executor.setQueueCapacity(1000);
    return executor;
}

常见问题问答(FAQ)

Q1:SSE和WebSocket选哪个?

A:单向后端推送选SSE(更简单、自动重连、兼容HTTP/2);双向交互或游戏场景必须用WebSocket。

Q2:SSE支持IE浏览器吗?

A:IE不支持原生EventSource,需引入eventsource-polyfill库,建议在项目兼容IE时使用WebSocket。

Q3:SSE连接意外断开,如何保证数据不丢失?

A:结合Redis持久化最近数据,断线重连后由前端发Last-Event-ID头,后端从断点续推。

Q4:高并发时每个用户一个SseEmitter会不会内存溢出?

A:会,必须设置超时时间、定期清理无用连接,并使用连接池,若单机承载超10万连接,需采用Nginx集群+Redis广播方案。

Q5:SSE可以传二进制数据吗?

A:原生SSE只支持UTF-8文本,若需二进制,可Base64编码传输,或直接改用WebSocket。


SSE是Java后端实现实时推送的“隐藏神器”,特别适合股票、监控、日志等单向数据流场景,通过Spring Boot + SseEmitter,你可以在几分钟内搭建一个生产级的实时推送服务,记住核心要点:心跳保活、超时清理、线程池隔离、水平扩展——这四点是保证SSE在高并发下稳定的关键,希望本文的案例和问答能帮你少走弯路,快速落地到你的项目中。

抱歉,评论功能暂时关闭!