发布于2026-07-14 阅读(0)
扫一扫,手机访问
在Web应用开发中,消息推送是个绕不开的话题——站内信、订单状态变更、系统告警通知,但凡涉及实时交互的场景,基本都离不开它。SSE(Server-Sent Events)作为一种轻量级的单向推送方案,优雅且高效,但一旦服务部署了多个实例,麻烦就来了:用户A连的是实例1,用户B连的是实例2,那A发出去的消息,怎么才能送到B手里?

上图把整个消息流转过程讲得很清楚:
- 建立连接:用户A和B分别连上不同的后端实例,每个实例各自维护着自己的SSE连接池。
- 发送消息:用户A发起私信请求,这个请求恰好落到了实例1上。
- Redis广播:实例1将消息发布到Redis的`station:message`频道,所有订阅了这个频道的实例都会收到。
- 推送消息:实例2发现目标用户B在自己身上,于是通过SSE连接把消息推过去。而实例1收到广播后一查,发现目标用户不在自己这,直接忽略。
先把依赖拉进来。Spring Boot Web本身对SSE的支持很到位,提供了`SseEmitter`这个核心类。`spring-boot-starter-data-redis`则提供了Redis连接和Pub/Sub能力。至于`commons-pool2`,它是连接池,生产环境必备,能避免频繁创建连接带来的性能损耗。
org.springframework.boot spring-boot-starter-web org.springframework.boot spring-boot-starter-data-redis org.apache.commons commons-pool2
`SseEmitterManager`是整个方案的主心骨,几个设计点值得细说:
@Component
@Slf4j
public class SseEmitterManager {
/**
* 用户ID -> (连接token -> SseEmitter)
* 一个用户可能有多个浏览器标签页,用token区分
*/
private final Map> userEmitters = new ConcurrentHashMap<>();
/**
* 建立SSE连接
* @param userId 用户ID
* @param token 连接标识(可用UUID)
*/
public SseEmitter connect(String userId, String token) {
// 超时时间设为0表示不超时,也可设置具体毫秒数
SseEmitter emitter = new SseEmitter(0L);
// 注册回调:连接关闭时清理
emitter.onCompletion(() -> removeEmitter(userId, token));
emitter.onTimeout(() -> removeEmitter(userId, token));
emitter.onError(e -> removeEmitter(userId, token));
// 存储连接
userEmitters.computeIfAbsent(userId, k -> new ConcurrentHashMap<>())
.put(token, emitter);
log.info("SSE connected: userId={}, token={}, total users={}",
userId, token, userEmitters.size());
return emitter;
}
/**
* 向指定用户推送消息
*/
public void sendToUser(String userId, String message) {
Map emitters = userEmitters.get(userId);
if (emitters == null || emitters.isEmpty()) {
log.debug("User {} not online, message stored for later", userId);
return;
}
// 向该用户所有连接推送
emitters.forEach((token, emitter) -> {
try {
emitter.send(SseEmitter.event()
.name("message")
.data(message));
} catch (IOException e) {
log.error("Send to user {} failed, removing emitter", userId, e);
removeEmitter(userId, token);
}
});
}
/**
* 全站广播
*/
public void broadcast(String message) {
userEmitters.forEach((userId, emitters) -> {
sendToUser(userId, message);
});
log.info("Broadcast message to {} users", userEmitters.size());
}
/**
* 获取当前在线人数
*/
public int getOnlineCount() {
return userEmitters.size();
}
private void removeEmitter(String userId, String token) {
Map emitters = userEmitters.get(userId);
if (emitters != null) {
emitters.remove(token);
if (emitters.isEmpty()) {
userEmitters.remove(userId);
}
}
}
}
`RedisMessageSubscriber`负责监听`station:message`频道。收到消息后,它做两件事:
这里有个容易被忽略的细节:每个后端实例都会收到自己发布的消息,所以推送前需要判断目标用户是否在当前实例上。这个判断逻辑其实隐含在`sendToUser`中——如果目标用户不在本实例的连接池里,就直接返回,不会报错。
@Component
@Slf4j
public class RedisMessageSubscriber implements MessageListener {
@Autowired
private SseEmitterManager sseEmitterManager;
@Autowired
private ObjectMapper objectMapper;
@Override
public void onMessage(Message message, byte[] pattern) {
try {
String channel = new String(message.getChannel());
String body = new String(message.getBody());
// 解析消息
StationMessage msg = objectMapper.readValue(body, StationMessage.class);
log.info("Received Redis message: channel={}, type={}, target={}",
channel, msg.getType(), msg.getTargetUserId());
// 根据消息类型分发
if ("user".equals(msg.getType())) {
// 私信:发给指定用户
sseEmitterManager.sendToUser(msg.getTargetUserId(), msg.getContent());
} else if ("broadcast".equals(msg.getType())) {
// 广播:发给所有在线用户
sseEmitterManager.broadcast(msg.getContent());
}
} catch (Exception e) {
log.error("Failed to process Redis message", e);
}
}
}
`RedisMessageListenerContainer`是Spring Data Redis提供的消息容器,它会自动管理订阅和监听线程。这里订阅的是`station:message`频道,你可以根据业务需要定义多个频道,比如`station:notice`、`station:system`。
序列化配置也很重要:key使用`StringRedisSerializer`保证可读性,value使用`Jackson2JsonRedisSerializer`来支持对象存储。
@Configuration
public class RedisConfig {
@Bean
public RedisMessageListenerContainer redisMessageListenerContainer(
RedisConnectionFactory connectionFactory,
RedisMessageSubscriber subscriber) {
RedisMessageListenerContainer container = new RedisMessageListenerContainer();
container.setConnectionFactory(connectionFactory);
// 订阅站内信频道
container.addMessageListener(subscriber, new ChannelTopic("station:message"));
return container;
}
@Bean
public RedisTemplate redisTemplate(RedisConnectionFactory factory) {
RedisTemplate template = new RedisTemplate<>();
template.setConnectionFactory(factory);
template.setKeySerializer(new StringRedisSerializer());
template.setValueSerializer(new Jackson2JsonRedisSerializer<>(Object.class));
return template;
}
}
`StationMessage`是消息的载体,在Redis中传输的JSON格式就对应这个结构。字段设计上:
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class StationMessage {
private String id; // 消息ID
private String type; // user:私信, broadcast:广播
private String targetUserId; // 私信时的目标用户ID
private String content; // 消息内容
private String senderId; // 发送者ID
private Long timestamp; // 时间戳
}
`MessageService`封装了发送逻辑,核心动作很简单:构造`StationMessage` -> 序列化为JSON -> `redisTemplate.convertAndSend()`。发布之后,所有实例的订阅者都会收到消息,相当于Redis帮我们做了一个“广播式”的跨实例通信。
这种设计的优点是:发送方不需要知道消息最终由哪个实例处理,也不需要维护实例之间的网络连接,所有协调工作都交给了Redis,省心又可靠。
@Service
@Slf4j
public class MessageService {
@Autowired
private RedisTemplate redisTemplate;
@Autowired
private SseEmitterManager sseEmitterManager;
@Autowired
private ObjectMapper objectMapper;
/**
* 发送私信
*/
public void sendPrivateMessage(String fromUserId, String toUserId, String content) {
StationMessage msg = StationMessage.builder()
.id(UUID.randomUUID().toString())
.type("user")
.targetUserId(toUserId)
.senderId(fromUserId)
.content(content)
.timestamp(System.currentTimeMillis())
.build();
try {
String json = objectMapper.writeValueAsString(msg);
// 发布到Redis,所有实例都会收到
redisTemplate.convertAndSend("station:message", json);
log.info("Private message sent: {} -> {}", fromUserId, toUserId);
} catch (JsonProcessingException e) {
log.error("Failed to serialize message", e);
}
}
/**
* 全站广播
*/
public void broadcast(String fromUserId, String content) {
StationMessage msg = StationMessage.builder()
.id(UUID.randomUUID().toString())
.type("broadcast")
.senderId(fromUserId)
.content(content)
.timestamp(System.currentTimeMillis())
.build();
try {
String json = objectMapper.writeValueAsString(msg);
redisTemplate.convertAndSend("station:message", json);
log.info("Broadcast message sent by: {}", fromUserId);
} catch (JsonProcessingException e) {
log.error("Failed to serialize broadcast", e);
}
}
}
Controller对外暴露了三个核心接口:
生产环境建议给这些接口加上认证鉴权,比如从token中解析userId,避免伪造身份。
@RestController
@RequestMapping("/api/sse")
@Slf4j
public class SseController {
@Autowired
private SseEmitterManager sseEmitterManager;
@Autowired
private MessageService messageService;
/**
* SSE连接端点
*/
@GetMapping(value = "/connect", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter connect(@RequestParam String userId) {
String token = UUID.randomUUID().toString();
return sseEmitterManager.connect(userId, token);
}
/**
* 发送私信
*/
@PostMapping("/private")
public ResponseEntity> sendPrivate(@RequestBody PrivateMessageRequest request) {
messageService.sendPrivateMessage(
request.getFromUserId(),
request.getToUserId(),
request.getContent()
);
return ResponseEntity.ok().build();
}
/**
* 全站广播
*/
@PostMapping("/broadcast")
public ResponseEntity> broadcast(@RequestBody BroadcastRequest request) {
messageService.broadcast(request.getFromUserId(), request.getContent());
return ResponseEntity.ok().build();
}
/**
* 获取在线人数
*/
@GetMapping("/online-count")
public ResponseEntity getOnlineCount() {
return ResponseEntity.ok(sseEmitterManager.getOnlineCount());
}
}
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
正版软件
正版软件
正版软件
正版软件
正版软件
1
2
3
7
8