如何在Spring Boot(Kotlin)中实现MongoDB ChangeStream并发送消息到RabbitMQ?
解决MongoDB ChangeStream类型不匹配问题并实现RabbitMQ消息推送
问题原因
你遇到的类型不匹配错误,是因为当前使用的MongoDB驱动返回的是异步反应式的ChangeStreamPublisher,但你声明的变量是同步的ChangeStreamIterable。这通常是因为依赖了异步驱动(mongodb-driver-reactivestreams),或者导入了错误的客户端类。
解决方案
1. 确保使用同步MongoDB驱动
如果你的项目是基于Spring Boot的同步架构(非WebFlux),请确保依赖的是同步驱动:
Gradle(Kotlin DSL)
dependencies { implementation("org.mongodb:mongodb-driver-sync:4.11.1") // 替换为对应版本 implementation("org.springframework.boot:spring-boot-starter-amqp") // RabbitMQ依赖 }
Maven
<dependencies> <dependency> <groupId>org.mongodb</groupId> <artifactId>mongodb-driver-sync</artifactId> <version>4.11.1</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> </dependencies>
2. 修正代码逻辑
- 不要把
main方法标记为@Bean,Spring Boot的@Bean需要返回可管理的Bean实例。 - 确保导入同步驱动的相关类,避免和异步类混淆。
- 补充RabbitMQ消息发送逻辑,处理增删改不同事件的字段提取。
修正后的代码示例:
import com.mongodb.client.MongoClients import com.mongodb.client.MongoCollection import com.mongodb.client.MongoDatabase import com.mongodb.client.model.Aggregates import com.mongodb.client.model.Filters import com.mongodb.client.model.FullDocument import org.bson.Document import org.slf4j.LoggerFactory import org.springframework.amqp.rabbit.core.RabbitTemplate import org.springframework.context.annotation.Bean import org.springframework.context.annotation.Configuration import java.util.concurrent.Executors @Configuration class VehicleChangeStreamConfig { private val log = LoggerFactory.getLogger(javaClass) @Bean fun vehicleChangeStreamListener(rabbitTemplate: RabbitTemplate): Runnable { return Runnable { log.info("车辆集合ChangeStream监听已启动") val uri = "mongodb://localhost:27017" MongoClients.create(uri).use { client -> val database: MongoDatabase = client.getDatabase("test") val collection: MongoCollection<Document> = database.getCollection("vehicles") val pipeline = listOf( Aggregates.match( Filters.`in`("operationType", listOf("insert", "update", "delete")) ) ) // 同步驱动下,watch返回ChangeStreamIterable,解决类型不匹配问题 val changeStream = collection.watch(pipeline) .fullDocument(FullDocument.UPDATE_LOOKUP) // 用单独线程执行监听,避免阻塞Spring启动流程 Executors.newSingleThreadExecutor().submit { changeStream.forEach { event -> log.info("收到集合变化事件: {}", event) // 提取id和market字段,区分不同操作类型处理 val vehicleId = event.getDocumentKey().getString("_id").value val market = when (event.operationTypeValue) { "delete" -> null // 删除事件无完整文档 else -> event.fullDocument?.getString("market")?.value } // 发送消息到RabbitMQ,假设队列名为"vehicle-changes" rabbitTemplate.convertAndSend("vehicle-changes", mapOf("id" to vehicleId, "market" to market)) } } } } } }
3. 启动监听
在Spring Boot启动类中,注入并启动监听任务:
import org.springframework.boot.autoconfigure.SpringBootApplication import org.springframework.boot.runApplication import org.springframework.context.ConfigurableApplicationContext @SpringBootApplication class VehicleApplication fun main(args: Array<String>) { val context = runApplication<VehicleApplication>(*args) val listener = context.getBean(Runnable::class.java) listener.run() }
关键说明
- 同步vs异步驱动:如果是WebFlux反应式架构,需使用
mongodb-driver-reactivestreams,并通过Flux.from(publisher).subscribe(...)的方式处理事件。 - 事件处理细节:删除事件的
fullDocument为空,需单独处理;更新事件通过FullDocument.UPDATE_LOOKUP可获取更新后的完整文档。 - 线程管理:ChangeStream的
forEach是阻塞调用,必须放在单独线程中执行,避免阻塞Spring Boot启动。
内容的提问来源于stack exchange,提问作者Max
相关产品推荐
相关产品推荐

