Java站内推送案例怎么实现

wen java案例 28

本文目录导读:

Java站内推送案例怎么实现

  1. 核心原理对比
  2. 案例一:基于 SSE (Server-Sent Events) 推送通知
  3. 案例二:基于 WebSocket (Spring Boot + WebSocket)
  4. 案例三:基于消息队列 (MQ) 的推送(高并发场景)
  5. 总结:如何选择?

Java 站内推送的实现方案有很多种,核心目标是服务端主动向浏览器/客户端发送消息,而不是客户端轮询。

以下是几种主流的 Java 实现方案,从简单到复杂,并附上对应的案例代码和适用场景。

核心原理对比

方案 技术栈 延迟 复杂度 浏览器兼容性 适用场景
短轮询 HTTP (Ajax) 高 (几秒) 最低 全兼容 实时性要求低,技术栈老旧
长轮询 HTTP (Comet) 全兼容 实时性要求一般,无法使用WebSocket
SSE (推荐) HTTP (Servlet) 低 (<1s) 现代浏览器 (不支持IE) 服务端单向推送 (如通知、股票)
WebSocket TCP (JSR 356) 极低 现代浏览器 双向通信 (如聊天、协同编辑)

基于 SSE (Server-Sent Events) 推送通知

这是目前 最推荐 的站内推送方式,服务端可以持续向客户端发送事件流。

特点:单向(服务端→客户端)、基于HTTP、轻量级、自动重连。

后端 Java 代码 (Spring Boot)

import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Sinks;
import java.time.Duration;
import java.time.LocalTime;
@RestController
public class NotificationController {
    // 1. 创建一个发布者(用于持续发送消息)
    private final Sinks.Many<String> sink = Sinks.many().multicast().onBackpressureBuffer();
    @GetMapping(value = "/subscribe", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public Flux<String> subscribe() {
        // 2. 返回一个Flux流,Spring会将其转换为SSE格式
        return sink.asFlux()
                .doOnCancel(() -> System.out.println("客户端断开连接"));
    }
    // 3. 模拟外部系统触发推送(比如管理员发通知)
    public void sendNotification(String message) {
        System.out.println("推送消息: " + message);
        sink.tryEmitNext("data: " + message);
    }
    // 4. 模拟周期性心跳
    @GetMapping("/tick")
    public Flux<String> tick() {
        return Flux.interval(Duration.ofSeconds(5))
                .map(i -> "data: 心跳 " + LocalTime.now().toString());
    }
}

前端 HTML/JS 代码

<!DOCTYPE html>
<html>
<head>站内推送 SSE</title>
</head>
<body>
<h1>站内通知</h1>
<div id="messages"></div>
<script>
    // 1. 创建EventSource连接到后端
    const eventSource = new EventSource('/subscribe');
    // 2. 监听消息事件
    eventSource.onmessage = function(event) {
        console.log('收到消息:', event.data);
        const div = document.getElementById('messages');
        div.innerHTML += `<p>${event.data}</p>`;
    };
    // 3. 监听错误/断开
    eventSource.onerror = function() {
        console.error('连接断开,浏览器会自动重连...');
    };
    // 4. 如果需要关闭连接
    // eventSource.close();
</script>
</body>
</html>

优点

  • 自动重连(断线后浏览器自动尝试重连)。
  • 兼容性好(现代浏览器都支持)。
  • 代码简单,Spring Boot 原生支持。

基于 WebSocket (Spring Boot + WebSocket)

适用于双向通信场景(如聊天、协同)。

WebSocket 配置

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.socket.server.standard.ServerEndpointExporter;
@Configuration
public class WebSocketConfig {
    @Bean
    public ServerEndpointExporter serverEndpointExporter() {
        return new ServerEndpointExporter();
    }
}

WebSocket 服务器端点

import javax.websocket.*;
import javax.websocket.server.ServerEndpoint;
import java.io.IOException;
import java.util.concurrent.CopyOnWriteArraySet;
@ServerEndpoint("/websocket/{userId}")
public class WebSocketServer {
    // 存储所有连接的客户端(用于广播)
    private static final CopyOnWriteArraySet<Session> sessions = new CopyOnWriteArraySet<>();
    private Session session;
    @OnOpen
    public void onOpen(Session session, @PathParam("userId") String userId) {
        this.session = session;
        sessions.add(session);
        System.out.println("用户 " + userId + " 已连接");
        // 可以保存 userId -> session 映射
    }
    @OnClose
    public void onClose() {
        sessions.remove(session);
        System.out.println("连接关闭");
    }
    @OnMessage
    public void onMessage(String message, Session session) {
        System.out.println("收到消息: " + message);
        // 处理业务逻辑...
    }
    @OnError
    public void onError(Session session, Throwable error) {
        error.printStackTrace();
    }
    // 广播给所有客户端
    public static void broadcast(String message) {
        for (Session s : sessions) {
            try {
                s.getBasicRemote().sendText(message);
            } catch (IOException e) {
                e.printStackTrace();
            }
        }
    }
}

前端代码

// 连接 WebSocket
const ws = new WebSocket('ws://localhost:8080/websocket/user123');
ws.onopen = function() {
    console.log('连接成功');
    ws.send('你好服务器');
};
ws.onmessage = function(event) {
    console.log('收到服务器消息:', event.data);
    // 更新页面...
};
ws.onclose = function() {
    console.log('连接关闭');
};

基于消息队列 (MQ) 的推送(高并发场景)

如果系统已经有了 RabbitMQ 或 Kafka,可以利用 MQ 广播 + WebSocket/SSE 的方式实现推送,这样做的优势是解耦和异步。

架构流程

  1. 业务系统发送消息到 MQ。
  2. 后端服务订阅 MQ 消息。
  3. 后端通过 WebSocketSSE 推送给对应的前端用户。
// 消费MQ消息
@RabbitListener(queues = "notification.queue")
public void handleNotification(NotificationMessage message) {
    // 根据userId找到对应的WebSocket Session (需要提前存储)
    WebSocketSession session = sessionMap.get(message.getUserId());
    if (session != null && session.isOpen()) {
        session.sendMessage(new TextMessage(message.getContent()));
    }
}

如何选择?

需求场景 推荐方案 理由
频率高、单向通知(如系统通知、更新公告) SSE 简单、自动重连、HTTP友好
聊天、协同编辑、需要双向发送 WebSocket 全双工、低延迟
简单、老旧系统、兼容IE 长轮询 (Long Polling) 兼容性最好
高并发、微服务架构、需要解耦 WebSocket + MQ 异步、削峰填谷
很少更新、精确数据(如订单状态) 接口 + 前端轮询 开发最快,后台只需写一个Controller

最推荐的组合SSE (服务端) + EventSource (前端) 对于绝大多数 Java 站内通知需求来说,是最快、最稳、代码最少的方案。

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