Java实现SSE案例

wen java案例 2

Java实现SSE案例:从零搭建服务端实时推送的完整指南

目录导读

  1. 什么是SSE?为何在Java中实现它?
  2. SSE vs WebSocket:选型对比与适用场景
  3. 环境准备与技术栈(JDK 17 + Spring Boot 3)
  4. 核心实现:基于Spring MVC的SSE端点
  5. 进阶玩法:异步线程池 + 心跳机制 + 断线重连
  6. 前端接入:EventSource API详解
  7. 实战案例:股票行情推送系统(附完整代码)
  8. 常见问题问答(FAQ)
  9. 性能调优与生产级注意事项
  10. 总结与最佳实践

什么是SSE?为何在Java中实现它?

SSE(Server-Sent Events,服务端发送事件)是一种基于HTTP的轻量级实时通信协议,它允许服务器主动向客户端推送数据,客户端只需通过EventSource接口即可订阅,与WebSocket不同,SSE是单向的(仅服务端→客户端),但它在Java中实现极其简单——你不需要额外的库,Spring框架原生支持。

Java实现SSE案例

为何选择Java实现?
Java生态成熟、Spring Boot的自动配置让SSE端点开发变得异常简洁,且天然支持线程池异步化,非常适合股票、日志流、AI聊天等实时场景。


SSE vs WebSocket:选型对比

维度 SSE WebSocket
通信方向 单向(服务端→客户端) 双向
协议 HTTP/1.1+(天然穿透代理) 独立协议(需要握手升级)
自动重连 内置(浏览器自动重连) 需手动实现
二进制支持 仅文本(Base64编码可绕过) 原生支持
实现复杂度 极低(Spring注解) 较高(需要处理连接状态)
适用场景 通知推送、实时图表、日志流 在线游戏、聊天、协同编辑

如果你的业务只需要服务端主动推送,SSE是更轻、更可靠的选择。


环境准备与技术栈

  • JDK 17+(推荐使用记录类型简化代码)
  • Spring Boot 3.x(内置Spring MVC 6)
  • Maven/Gradle(构建工具)
  • Lombok(可选,减少样板代码)
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-web</artifactId>
</dependency>

核心实现:基于Spring MVC的SSE端点

基础版:直接返回SseEmitter

@RestController
@RequestMapping("/api/sse")
public class SseController {
    @GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public SseEmitter stream() {
        SseEmitter emitter = new SseEmitter(0L); // 0表示永不超时
        // 模拟数据推送
        new Thread(() -> {
            try {
                for (int i = 0; i < 10; i++) {
                    emitter.send(SseEmitter.event()
                            .name("message")
                            .data(Map.of("id", i, "time", LocalTime.now().toString())));
                    Thread.sleep(1000);
                }
                emitter.complete();
            } catch (Exception e) {
                emitter.completeWithError(e);
            }
        }).start();
        return emitter;
    }
}

关键点

  • produces = MediaType.TEXT_EVENT_STREAM_VALUE 指定响应类型为SSE流。
  • SseEmitter.event().name("xxx") 自定义事件名,前端监听对应事件。
  • 必须手动调用complete()completeWithError()来结束连接。

改进版:使用ExecutorService线程池

private final ExecutorService executor = Executors.newCachedThreadPool();
@GetMapping("/stream2")
public SseEmitter stream2() {
    SseEmitter emitter = new SseEmitter(60_000L); // 60秒超时
    executor.execute(() -> {
        try {
            // 业务逻辑
            emitter.send("初始化数据...");
            emitter.complete();
        } catch (IOException e) {
            emitter.completeWithError(e);
        }
    });
    emitter.onTimeout(() -> System.out.println("连接超时"));
    emitter.onCompletion(() -> System.out.println("连接关闭"));
    return emitter;
}

进阶玩法:心跳机制与断线重连

在生产环境中,代理服务器(如Nginx)默认会关闭空闲连接,必须定期发送注释(如heartbeat)保持活跃。

@Scheduled(fixedRate = 15000)
public void heartbeat() {
    // 遍历所有活跃的SseEmitter,发送注释
    activeEmitters.forEach((id, emitter) -> {
        try {
            emitter.send(SseEmitter.event().comment("heartbeat"));
        } catch (Exception e) {
            activeEmitters.remove(id);
        }
    });
}

前端自动重连EventSource内建重连机制,服务端只需正常发送数据,客户端断开后会自动尝试重新连接(默认延迟3秒)。


前端接入:EventSource API详解

const eventSource = new EventSource('http://localhost:8080/api/sse/stream');
// 监听默认消息(无event字段)
eventSource.onmessage = (event) => {
    console.log('收到数据:', JSON.parse(event.data));
};
// 监听自定义事件
eventSource.addEventListener('message', (event) => {
    // 如果服务端设置了事件名
});
// 连接错误处理
eventSource.onerror = (event) => {
    if (event.target.readyState === EventSource.CLOSED) {
        console.log('连接已关闭');
    } else {
        console.log('正在重连...');
    }
};

实战案例:股票行情推送系统(完整代码)

@RestController
@RequestMapping("/stock")
public class StockController {
    private final Map<String, SseEmitter> clients = new ConcurrentHashMap<>();
    @GetMapping("/subscribe/{code}")
    public SseEmitter subscribe(@PathVariable String code) {
        SseEmitter emitter = new SseEmitter(0L);
        clients.put(code + "-" + UUID.randomUUID(), emitter);
        emitter.onCompletion(() -> clients.remove(code));
        emitter.onTimeout(() -> clients.remove(code));
        return emitter;
    }
    @PostMapping("/push/{code}")
    public String pushPrice(@PathVariable String code, @RequestBody PriceData data) {
        for (Map.Entry<String, SseEmitter> entry : clients.entrySet()) {
            if (entry.getKey().startsWith(code)) {
                try {
                    entry.getValue().send(SseEmitter.event()
                            .name("price")
                            .data(data));
                } catch (IOException e) {
                    clients.remove(entry.getKey());
                }
            }
        }
        return "推送成功";
    }
}

前端调用

const es = new EventSource('/stock/subscribe/600519');
es.addEventListener('price', (e) => {
    document.getElementById('price').textContent = JSON.parse(e.data).price;
});

常见问题问答(FAQ)

Q1: SSE连接为什么会频繁断开?
A: 很可能是代理服务器空闲超时所致,解决方案:① 服务端开启心跳线程定时发送注释;② 配置Nginx的proxy_read_timeout

Q2: 多个客户端共享同一个Emitter会怎样?
A: 每个客户端必须有独立的SseEmitter实例,否则数据会串流,务必使用Map存储每个连接的Emitter。

Q3: SSE支持自定义HTTP状态码吗?
A: 支持,但一旦响应头已发送则无法修改,建议在业务逻辑开始前设置响应状态。

Q4: 如何处理客户端异常断开?
A: 在onCompletiononTimeout回调中清理资源,同时捕获IOException


性能调优与生产级注意事项

  1. 连接限制:默认Tomcat支持200个并发线程,若推送频率高需调大server.tomcat.threads.max
  2. 异步化:推送业务务必使用@Async或独立线程池,避免阻塞Tomcat工作线程。
  3. 数据压缩:对大型JSON启用GZIP压缩(server.compression.enabled=true)。
  4. 监控:通过Metrics记录活跃连接数、推送失败率。
  5. 安全:生产环境使用HTTPS并添加鉴权(如JWT放在URL参数,SSE不支持自定义Header)。

总结与最佳实践

SSE是Java世界中最被低估的实时通信方案,通过Spring Boot的SseEmitter,你可以在30分钟内完成一个可靠的实时推送系统,最佳实践归纳:

  • ✅ 优先选择SSE而非WebSocket,除非需要双向通信
  • ✅ 务必实现心跳机制应对网络代理
  • ✅ 使用线程池隔离业务逻辑,保证响应速度
  • ✅ 前端配合EventSource自动重连,提升用户体验

拿起你的IDE,动手实现第一个Java SSE案例吧!如果你有任何问题,欢迎在评论区留言讨论。

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