基于Spring Reactive与Server Side Events的实时数据推送方案问询
基于Spring Reactive的实时数据推送解决方案
听起来你之前用capped collection的方案确实碰到了它的核心瓶颈——没法修改元素大小,这确实挺头疼的。不过别担心,基于Spring Reactive + MongoDB,我们有更灵活的方案能完美解决你的需求:利用MongoDB Change Streams结合SSE/WebSocket,完全避开capped集合的限制,实现数据变更后的实时推送。
核心思路:用MongoDB Change Streams替代Tailable游标
MongoDB的Change Streams是专门为监听集合CRUD变更设计的特性,它支持普通集合(不需要capped),能捕捉insert、update、delete等所有类型的操作,而且完全兼容Spring Reactive的响应式编程模型,完美适配你的技术栈。
前置准备:启用MongoDB Change Streams
首先你的MongoDB需要运行在副本集或分片集群模式下(单节点默认不支持,本地开发可以快速搭建一个单节点副本集,网上有很多教程,几步就能搞定)。
方案1:SSE(Server-Sent Events)单向推送
SSE是基于HTTP的单向推送机制,前端用EventSource就能轻松接收,非常适合你这种"后端主动推送更新后全量列表"的场景。
第一步:实现Change Streams监听服务
@Service public class ObjectChangeMonitor { private final ReactiveMongoTemplate reactiveMongoTemplate; private final SsePushService ssePushService; // 构造函数注入依赖 public ObjectChangeMonitor(ReactiveMongoTemplate reactiveMongoTemplate, SsePushService ssePushService) { this.reactiveMongoTemplate = reactiveMongoTemplate; this.ssePushService = ssePushService; } @PostConstruct public void startMonitoring() { // 监听指定集合的所有变更事件 reactiveMongoTemplate.changeStream(ObjectEntity.class) .watchCollection("your-target-collection") .listen() // 每次变更后查询最新的全量对象列表 .flatMap(changeEvent -> reactiveMongoTemplate.findAll(ObjectEntity.class).collectList()) // 推送给所有在线前端客户端 .subscribe(objectList -> ssePushService.broadcastUpdate(objectList)); } }
第二步:实现SSE广播服务
@Service public class SsePushService { private final List<SseEmitter> activeEmitters = new CopyOnWriteArrayList<>(); private final ReactiveMongoTemplate reactiveMongoTemplate; public SsePushService(ReactiveMongoTemplate reactiveMongoTemplate) { this.reactiveMongoTemplate = reactiveMongoTemplate; } @GetMapping("/stream/object-updates") public SseEmitter establishConnection() { SseEmitter emitter = new SseEmitter(Long.MAX_VALUE); activeEmitters.add(emitter); // 客户端断开连接时自动移除emitter emitter.onCompletion(() -> activeEmitters.remove(emitter)); emitter.onTimeout(() -> activeEmitters.remove(emitter)); // 首次连接时发送当前全量列表 reactiveMongoTemplate.findAll(ObjectEntity.class) .collectList() .subscribe(list -> { try { emitter.send(SseEmitter.event().data(list)); } catch (IOException e) { emitter.completeWithError(e); } }); return emitter; } public void broadcastUpdate(List<ObjectEntity> updatedList) { // 遍历所有活跃连接推送更新,移除失效的连接 activeEmitters.removeIf(emitter -> { try { emitter.send(SseEmitter.event().data(updatedList)); return false; } catch (IOException e) { emitter.completeWithError(e); return true; } }); } }
方案2:WebSocket双向推送
如果之后需要扩展双向通信(比如前端可以发送指令给后端),WebSocket是更合适的选择,Spring WebFlux原生支持WebSocket响应式处理。
第一步:实现WebSocket处理器
@Component public class ObjectWebSocketHandler implements WebSocketHandler { private final ReactiveMongoTemplate reactiveMongoTemplate; private final Flux<List<ObjectEntity>> updateStream; public ObjectWebSocketHandler(ReactiveMongoTemplate reactiveMongoTemplate) { this.reactiveMongoTemplate = reactiveMongoTemplate; // 创建可共享的变更事件流,所有客户端共享同一个流 ConnectableFlux<List<ObjectEntity>> connectableFlux = reactiveMongoTemplate.changeStream(ObjectEntity.class) .watchCollection("your-target-collection") .flatMap(ignored -> reactiveMongoTemplate.findAll(ObjectEntity.class).collectList()) .publish(); this.updateStream = connectableFlux.autoConnect(); } @Override public Mono<Void> handle(WebSocketSession session) { // 向客户端推送更新列表,同时处理客户端消息(这里留空,只做单向推送) return session.send(updateStream.map(session::textMessage)) .and(session.receive().doOnNext(msg -> {}).then()); } }
第二步:配置WebSocket路由
@Configuration public class WebSocketConfig { @Bean public RouterFunction<ServerResponse> webSocketRoute(ObjectWebSocketHandler handler) { return RouterFunctions.route( RequestPredicates.GET("/ws/object-updates"), request -> ServerResponse.ok().build(session -> handler.handle(session)) ); } }
优化建议:增量更新代替全量查询
上面的示例每次变更都查询全量列表,适合数据量不大的场景。如果数据量较大,你可以根据Change Stream的事件类型做增量更新,减少数据库开销:
reactiveMongoTemplate.changeStream(ObjectEntity.class) .watchCollection("your-target-collection") .listen() .subscribe(changeEvent -> { ChangeStreamEvent.Body<ObjectEntity> eventBody = changeEvent.getBody(); switch (eventBody.getOperationType()) { case INSERT: // 处理新增:向客户端列表添加新对象 break; case UPDATE: // 处理更新:替换客户端列表中对应ID的对象 break; case DELETE: // 处理删除:从客户端列表移除对应ID的对象 break; default: break; } });
方案核心优势
- 完全摆脱capped collection的限制:可以自由修改对象大小、删除对象,和普通集合用法完全一致
- 支持全类型CRUD事件监听:新增、修改、删除都能精准捕捉
- 完美适配Spring Reactive模型:所有操作都是非阻塞的,符合你的技术栈要求
内容的提问来源于stack exchange,提问作者minihulk22
相关产品推荐
相关产品推荐

