Java消息推送案例如何编写

wen java案例 28

Java消息推送案例如何编写:从零构建实时通信系统

目录导读

  1. 什么是消息推送?为什么需要它?
  2. 消息推送的核心技术选型
  3. 基于WebSocket的Java推送实战案例
  4. 基于SSE(Server-Sent Events)的轻量级推送
  5. 结合消息中间件(如RabbitMQ)的推送架构
  6. 常见问题与问答(Q&A)
  7. 性能优化与生产环境注意事项

什么是消息推送?为什么需要它?

消息推送,顾名思义,是指服务器主动向客户端发送数据,而不是等待客户端轮询请求,在现代Web应用、移动App、物联网系统中,消息推送是实现实时交互(如聊天、通知、股票行情、在线编辑)的核心技术。

Java消息推送案例如何编写

传统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 → 其他服务

核心思路

  1. 每个WebSocket服务器实例订阅同一个RabbitMQ队列。
  2. 业务服务将推送消息发送到RabbitMQ Exchange。
  3. 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(),并在服务端设置合理的超时时间,可以使用SseEmitteronCompletiononTimeout回调清理资源。

Q3: 消息推送丢失怎么办?
A: WebSocket默认不保证消息可靠投递,生产环境应结合消息队列(如RabbitMQ)实现ACK机制,并在客户端记录消息序列号,发现丢失后请求重推。

Q4: 如何实现群组/多用户推送?
A: 在WebSocket连接时,将用户ID及其所在群组关系存入Redis,推送时通过群组ID查询所有在线用户ID,然后逐个调用sendToUser

Q5: 生产环境如何监控推送成功率?
A: 在推送接口中加入埋点(如发送计数、成功计数),使用Prometheus + Grafana监控,同时记录推送日志到ELK,方便排查。


性能优化与生产环境注意事项

  1. 限制最大连接数:防止内存溢出,在WebSocket握手前,可以拦截器判断当前连接数是否超过阈值。
  2. 使用NIO框架:WebSocket在底层使用Netty时性能更优,Spring Boot内置的Tomcat也支持NIO,但若并发上万,建议用Netty替换。
  3. 消息压缩:推送大量JSON数据时,启用GZIP压缩(需服务端和客户端都支持)。
  4. 后端异步化:推送任务不要阻塞业务线程,使用线程池或响应式WebFlux。
  5. 集群节点间通知:当某个WebSocket实例与客户端建立连接后,如果另一个实例需要推送给同一用户,需通过Redis Pub/Sub或消息队列广播用户连接信息。

最终总结:Java消息推送的核心在于“选对技术、管好连接、保证可靠”,WebSocket适合双向实时交互,SSE适合单向轻量推送,而消息中间件则是高并发场景下的“定海神针”,无论选择哪种方式,务必关注连接管理、心跳保活和消息可靠性这三个关键点。


注:文中示例代码基于Spring Boot 2.x,若使用Spring Boot 3.x(基于Jakarta EE),配置类需替换为jakarta.websocket,但核心逻辑一致。

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