如何通过REST端点临时暂停并恢复MongoDB Change Stream?
实现带启停与恢复功能的MongoDB变更流REST端点
核心思路
要实现变更流的暂停和基于Resume Token的恢复,关键在于:
- 存储变更流的Resume Token(记录流的最后处理位置)
- 管控变更流的生命周期,支持启动与停止操作
- 启动时若存在有效Resume Token,基于它从断点恢复流
代码实现
1. 变更流服务类(核心逻辑)
创建单例服务类,封装变更流的启停、Resume Token管理逻辑:
import org.springframework.data.mongodb.core.ReactiveMongoTemplate; import org.springframework.data.mongodb.core.changeStream.ChangeStreamEvent; import org.springframework.data.mongodb.core.changeStream.ChangeStreamOptions; import org.springframework.stereotype.Service; import reactor.core.publisher.Flux; import reactor.core.publisher.Sinks; import reactor.core.Disposable; import java.util.concurrent.atomic.AtomicReference; @Service public class ExampleChangeStreamService { private final ReactiveMongoTemplate reactiveMongoTemplate; // 线程安全存储Resume Token private final AtomicReference<Object> resumeToken = new AtomicReference<>(); // 统一管理变更流输出,适配外部订阅 private final Sinks.Many<Example> eventSink = Sinks.many().multicast().onBackpressureBuffer(); // 记录当前流的订阅句柄,用于停止流 private volatile Disposable currentStreamDisposable; public ExampleChangeStreamService(ReactiveMongoTemplate reactiveMongoTemplate) { this.reactiveMongoTemplate = reactiveMongoTemplate; } // 提供外部订阅变更流的入口 public Flux<Example> getChangeStreamEvents() { return eventSink.asFlux(); } // 启动/恢复变更流 public void startStream() { if (currentStreamDisposable != null && !currentStreamDisposable.isDisposed()) { return; // 流已在运行,直接返回 } ChangeStreamOptions.ChangeStreamOptionsBuilder optionsBuilder = ChangeStreamOptions.builder() .returnFullDocumentOnUpdate(); // 存在有效Resume Token时,设置恢复起点 Object token = resumeToken.get(); if (token != null) { optionsBuilder.resumeToken(token); } ChangeStreamOptions options = optionsBuilder.build(); // 创建变更流并订阅,将事件转发到Sink currentStreamDisposable = reactiveMongoTemplate.changeStream("collection", options, Example.class) .filter(event -> event.getOperationType() != null) .doOnNext(event -> resumeToken.set(event.getResumeToken())) // 实时更新Resume Token .mapNotNull(ChangeStreamEvent::getBody) .subscribe(eventSink::tryEmitNext, eventSink::tryEmitError); } // 停止变更流 public void stopStream() { if (currentStreamDisposable != null && !currentStreamDisposable.isDisposed()) { currentStreamDisposable.dispose(); currentStreamDisposable = null; } } }
2. REST控制器(暴露端点)
创建REST控制器,提供启停变更流的接口:
import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; import reactor.core.publisher.Mono; import reactor.core.publisher.Flux; @RestController public class ChangeStreamController { private final ExampleChangeStreamService changeStreamService; public ChangeStreamController(ExampleChangeStreamService changeStreamService) { this.changeStreamService = changeStreamService; } // 启动/恢复变更流 @PostMapping("/change-stream/start") public Mono<String> startChangeStream() { changeStreamService.startStream(); return Mono.just("变更流已启动/恢复"); } // 停止变更流 @PostMapping("/change-stream/stop") public Mono<String> stopChangeStream() { changeStreamService.stopStream(); return Mono.just("变更流已停止"); } // 可选:提供订阅变更流事件的端点 @GetMapping("/change-stream/events") public Flux<Example> streamEvents() { return changeStreamService.getChangeStreamEvents(); } }
关键细节说明
- Resume Token持久化:示例用
AtomicReference做内存存储,若需服务重启后仍能恢复,可将Token存入Redis、数据库等持久化介质 - Sinks的作用:通过
Sinks.Many统一管理事件输出,避免每次启停都创建新Flux,保证外部订阅者能持续接收事件 - 线程安全:用
volatile和AtomicReference保证多线程环境下的状态安全,避免并发启停引发异常 - 恢复逻辑:启动时检测到Resume Token,会通过
resumeToken()方法传入配置,MongoDB将从该断点继续推送后续事件
内容的提问来源于stack exchange,提问作者junelane
相关产品推荐
相关产品推荐

