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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 04:02:47