You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.14 08:15:55