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

如何在Spring Boot中启动MongoDB Change Stream监听程序?

在Spring Boot中启动MongoDB Change Stream监听的方法

首先优化你的现有代码,Spring Boot中建议使用自动配置的Mongo资源,不要手动创建MongoClient,避免重复实例和资源泄漏问题:

import com.mongodb.client.*;
import org.bson.Document;
import org.bson.conversions.Bson;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.mongodb.core.MongoTemplate;
import org.springframework.stereotype.Component;

import java.util.Arrays;
import java.util.List;

@Component
public class MongoDBChangeStream {

    private final MongoCollection<Document> teamsCollection;

    // 注入Spring自动配置的MongoTemplate,通过它获取集合
    @Autowired
    public MongoDBChangeStream(MongoTemplate mongoTemplate) {
        this.teamsCollection = mongoTemplate.getDb().getCollection("teams");
    }

    public void startChangeStreamListener() {
        // 创建操作类型过滤管道
        List<Bson> pipeline = Arrays.asList(
                Aggregates.match(
                        Filters.in("operationType", Arrays.asList("insert", "update", "delete"))
                )
        );

        // 创建Change Stream,UPDATE_LOOKUP会返回更新后的完整文档
        ChangeStreamIterable<Document> changeStream = teamsCollection.watch(pipeline)
                .fullDocument(FullDocument.UPDATE_LOOKUP);

        // 迭代监听变更——注意这个循环是阻塞的,必须放在独立线程中执行
        try (MongoCursor<ChangeStreamDocument<Document>> cursor = changeStream.cursor()) {
            while (cursor.hasNext()) {
                ChangeStreamDocument<Document> changeEvent = cursor.next();
                switch (changeEvent.getOperationType()) {
                    case INSERT:
                        System.out.println("MongoDB Change Stream检测到插入操作: " + changeEvent.getFullDocument());
                        break;
                    case UPDATE:
                        System.out.println("MongoDB Change Stream检测到更新操作: " + changeEvent.getFullDocument());
                        break;
                    case DELETE:
                        System.out.println("MongoDB Change Stream检测到删除操作: " + changeEvent.getDocumentKey());
                        break;
                }
            }
        } catch (Exception e) {
            System.err.println("Change Stream监听异常: " + e.getMessage());
            // 这里可以根据需求添加重连逻辑,避免单次异常导致监听终止
        }
    }
}

接下来提供几种可靠的启动方式,核心是把阻塞的监听逻辑放到异步线程中,避免阻塞Spring Boot启动流程:

方式1:用@PostConstruct配合线程池

创建一个启动器组件,在Spring Bean初始化完成后启动监听:

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;

import javax.annotation.PostConstruct;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

@Component
public class ChangeStreamStarter {

    private final MongoDBChangeStream changeStream;
    private ExecutorService executor;

    @Autowired
    public ChangeStreamStarter(MongoDBChangeStream changeStream) {
        this.changeStream = changeStream;
    }

    @PostConstruct
    public void startListener() {
        // 创建单线程池专门运行监听任务
        executor = Executors.newSingleThreadExecutor();
        executor.submit(changeStream::startChangeStreamListener);
    }

    // 应用关闭时优雅关闭线程池
    @Override
    public void destroy() throws Exception {
        if (executor != null) {
            executor.shutdown();
        }
    }
}

方式2:实现ApplicationRunner接口

在Spring Boot启动完成后执行监听任务,同样要放到异步线程:

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.stereotype.Component;

import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

@Component
public class ChangeStreamRunner implements ApplicationRunner {

    private final MongoDBChangeStream changeStream;

    @Autowired
    public ChangeStreamRunner(MongoDBChangeStream changeStream) {
        this.changeStream = changeStream;
    }

    @Override
    public void run(ApplicationArguments args) {
        ExecutorService executor = Executors.newSingleThreadExecutor();
        executor.submit(changeStream::startChangeStreamListener);
    }
}

方式3:使用Spring@Async注解

  1. 先在启动类上开启异步支持:
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.scheduling.annotation.EnableAsync;

@SpringBootApplication
@EnableAsync
public class YourApplication {
    public static void main(String[] args) {
        SpringApplication.run(YourApplication.class, args);
    }
}
  1. 给监听方法添加@Async注解:
// ... 其他代码不变

@Async
public void startChangeStreamListener() {
    // 原有的监听逻辑
}
  1. 创建初始化组件,在应用就绪后调用监听方法:
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Component;

@Component
public class ChangeStreamInitializer {

    private final MongoDBChangeStream changeStream;

    @Autowired
    public ChangeStreamInitializer(MongoDBChangeStream changeStream) {
        this.changeStream = changeStream;
    }

    @EventListener(ApplicationReadyEvent.class)
    public void init() {
        changeStream.startChangeStreamListener();
    }
}

关键注意事项

  • 线程隔离:绝对不能在主线程或Spring初始化线程中直接运行阻塞的Change Stream循环,否则会导致应用启动卡住。
  • 异常处理:必须给监听循环添加异常捕获,防止因数据库断开、网络波动等异常导致监听线程直接终止。
  • 资源清理:记得在应用关闭时关闭线程池和Mongo资源,避免内存泄漏。

内容的提问来源于stack exchange,提问作者TheStranger

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 05:01:08