如何在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注解
- 先在启动类上开启异步支持:
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); } }
- 给监听方法添加@Async注解:
// ... 其他代码不变 @Async public void startChangeStreamListener() { // 原有的监听逻辑 }
- 创建初始化组件,在应用就绪后调用监听方法:
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
相关产品推荐
相关产品推荐

