Java消息推送案例如何编写:从零构建实时通信系统
目录导读
- 什么是消息推送?为什么需要它?
- 消息推送的核心技术选型
- 基于WebSocket的Java推送实战案例
- 基于SSE(Server-Sent Events)的轻量级推送
- 结合消息中间件(如RabbitMQ)的推送架构
- 常见问题与问答(Q&A)
- 性能优化与生产环境注意事项
什么是消息推送?为什么需要它?
消息推送,顾名思义,是指服务器主动向客户端发送数据,而不是等待客户端轮询请求,在现代Web应用、移动App、物联网系统中,消息推送是实现实时交互(如聊天、通知、股票行情、在线编辑)的核心技术。

传统HTTP轮询的痛点:客户端每隔几秒定时发送请求,服务器无论数据有无变化都返回响应,这不仅浪费带宽,还会造成服务器压力剧增,且延迟较高。
推送的典型场景:
- 即时通讯(IM)
- 系统告警通知
- 数据看板实时刷新
- 在线协作编辑
- 电商订单状态更新
消息推送的核心技术选型
Java生态中,常见的消息推送技术有三种,各有优劣:
| 技术 | 协议 | 适用场景 | 延迟 | 复杂度 |
|---|---|---|---|---|
| WebSocket | 全双工 | 高频率实时交互 | 极低 | 中 |
| SSE (Server-Sent Events) | 单向(服务端→客户端) | 实时通知、日志流 | 低 | 低 |
| 轮询/长轮询 | HTTP | 兼容性要求高、简单场景 | 高 | 低 |
推荐选择:如果只需要服务器向客户端推送(如通知、行情),SSE更轻量;如果需要双向通信(如聊天),WebSocket是首选,对于高并发、高可靠性场景,建议在WebSocket底层引入消息中间件。
基于WebSocket的Java推送实战案例
1 环境准备
- JDK 8+
- Spring Boot 2.x 或 3.x
- 依赖(Maven):
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-websocket</artifactId> </dependency>
2 WebSocket服务端核心代码
// 1. 配置类:启用WebSocket
@Configuration
@EnableWebSocket
public class WebSocketConfig implements WebSocketConfigurer {
@Override
public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
registry.addHandler(new MyWebSocketHandler(), "/ws/push")
.setAllowedOrigins("*"); // 生产环境需限制域名
}
}
// 2. 处理器:管理连接与消息广播
public class MyWebSocketHandler extends TextWebSocketHandler {
// 存储所有在线客户端(ConcurrentHashMap保证线程安全)
private static final Map<String, WebSocketSession> clients = new ConcurrentHashMap<>();
@Override
public void afterConnectionEstablished(WebSocketSession session) {
String userId = session.getAttributes().get("userId").toString();
clients.put(userId, session);
System.out.println("用户 " + userId + " 已连接");
}
@Override
protected void handleTextMessage(WebSocketSession session, TextMessage message) {
// 接收客户端消息,可做业务处理
String payload = message.getPayload();
// 示例:广播消息给所有用户
broadcast("用户 " + session.getId() + " 说: " + payload);
}
@Override
public void afterConnectionClosed(WebSocketSession session, CloseStatus status) {
clients.values().remove(session);
}
// 封装推送方法:推送给指定用户
public static void sendToUser(String userId, String message) {
WebSocketSession session = clients.get(userId);
if (session != null && session.isOpen()) {
session.sendMessage(new TextMessage(message));
}
}
// 封装广播方法
public static void broadcast(String message) {
clients.values().forEach(session -> {
try {
session.sendMessage(new TextMessage(message));
} catch (IOException e) {
e.printStackTrace();
}
});
}
}
// 3. 业务中调用推送(例如Controller)
@RestController
public class PushController {
@PostMapping("/push/toUser")
public String pushToUser(@RequestParam String userId, @RequestParam String msg) {
MyWebSocketHandler.sendToUser(userId, msg);
return "推送成功";
}
}
3 前端JavaScript连接示例
let socket = new WebSocket("ws://localhost:8080/ws/push");
socket.onopen = function() {
console.log("连接已建立");
socket.send("Hello Server");
};
socket.onmessage = function(event) {
console.log("收到推送:", event.data);
};
基于SSE(Server-Sent Events)的轻量级推送
SSE只需要一个HTTP长连接,服务端可以持续推送文本数据,且浏览器原生支持(EventSource API)。
1 Spring Boot SSE服务端代码
@RestController
public class SSEPushController {
@GetMapping(value = "/sse/push", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter push() {
SseEmitter emitter = new SseEmitter(0L); // 0表示超时时间无限
// 模拟推送任务
new Thread(() -> {
try {
for (int i = 0; i < 10; i++) {
emitter.send(SseEmitter.event()
.name("message")
.data("第" + i + "条推送"));
Thread.sleep(2000);
}
emitter.complete();
} catch (Exception e) {
emitter.completeWithError(e);
}
}).start();
return emitter;
}
}
2 前端SSE客户端
let source = new EventSource("/sse/push");
source.addEventListener("message", function(event) {
console.log("SSE推送:", event.data);
});
SSE与WebSocket对比:
- SSE更简单,但仅支持文本推送(非二进制)。
- WebSocket可双向通信,但需要更复杂的握手和心跳保活。
结合消息中间件(如RabbitMQ)的推送架构
当系统有多个微服务实例,或需要保证消息不丢失、支持高并发时,建议引入消息中间件:
架构:客户端 → WebSocket集群 → RabbitMQ → 其他服务
核心思路:
- 每个WebSocket服务器实例订阅同一个RabbitMQ队列。
- 业务服务将推送消息发送到RabbitMQ Exchange。
- RabbitMQ将消息路由到所有订阅的队列,WebSocket集群消费并推送给对应客户端。
示例代码片段(RabbitMQ + WebSocket):
// 消息监听器(在WebSocket服务中)
@Component
public class PushMessageListener {
@RabbitListener(queues = "push.queue")
public void handleMessage(String msg) {
// 解析消息,获取目标userId,然后调用WebSocket发送
MyWebSocketHandler.sendToUser(userId, msg);
}
}
常见问题与问答(Q&A)
Q1: WebSocket连接为什么经常断开?如何处理心跳?
A: 网络中间件(如Nginx、防火墙)可能会断开空闲连接,解决方案:在WebSocket处理器中实现定时心跳(例如每隔30秒发送Ping消息),并在客户端监听pong事件。
Q2: 使用SSE时,客户端连接数过多导致内存泄漏怎么办?
A: 务必在每个SseEmitter完成或出错后调用emitter.complete(),并在服务端设置合理的超时时间,可以使用SseEmitter的onCompletion和onTimeout回调清理资源。
Q3: 消息推送丢失怎么办?
A: WebSocket默认不保证消息可靠投递,生产环境应结合消息队列(如RabbitMQ)实现ACK机制,并在客户端记录消息序列号,发现丢失后请求重推。
Q4: 如何实现群组/多用户推送?
A: 在WebSocket连接时,将用户ID及其所在群组关系存入Redis,推送时通过群组ID查询所有在线用户ID,然后逐个调用sendToUser。
Q5: 生产环境如何监控推送成功率?
A: 在推送接口中加入埋点(如发送计数、成功计数),使用Prometheus + Grafana监控,同时记录推送日志到ELK,方便排查。
性能优化与生产环境注意事项
- 限制最大连接数:防止内存溢出,在WebSocket握手前,可以拦截器判断当前连接数是否超过阈值。
- 使用NIO框架:WebSocket在底层使用Netty时性能更优,Spring Boot内置的Tomcat也支持NIO,但若并发上万,建议用Netty替换。
- 消息压缩:推送大量JSON数据时,启用GZIP压缩(需服务端和客户端都支持)。
- 后端异步化:推送任务不要阻塞业务线程,使用线程池或响应式WebFlux。
- 集群节点间通知:当某个WebSocket实例与客户端建立连接后,如果另一个实例需要推送给同一用户,需通过Redis Pub/Sub或消息队列广播用户连接信息。
最终总结:Java消息推送的核心在于“选对技术、管好连接、保证可靠”,WebSocket适合双向实时交互,SSE适合单向轻量推送,而消息中间件则是高并发场景下的“定海神针”,无论选择哪种方式,务必关注连接管理、心跳保活和消息可靠性这三个关键点。
注:文中示例代码基于Spring Boot 2.x,若使用Spring Boot 3.x(基于Jakarta EE),配置类需替换为jakarta.websocket,但核心逻辑一致。