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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 23:15:32