【问题标题】:Transaction synchronisation: how to create ChainedKafkaTransactionManager bean with Reactor Kafka and R2DBC事务同步:如何使用 Reactor Kafka 和 R2DBC 创建 ChainedKafkaTransactionManager bean
【发布时间】:2021-04-11 06:19:45
【问题描述】:

我的 Spring Boot (WebFlux/R2DBC/Reactor Kafka) 应用程序中有以下使用者

    @EventListener(ApplicationStartedEvent::class)
    fun onMyEvent() {
        kafkaReceiver
            .receive()
            .doOnNext { record ->
                val myEvent = record.value()
                myService.deleteSomethingFromDbById(myEvent.myId)
                .thenEmpty {
                    record.receiverOffset().acknowledge()
                }.subscribe()
            }
            .subscribe()
    }

我想为 Kafka 和 DB 事务添加事务同步。阅读文档和一些 stakoverflow 问题后

似乎ChainedKafkaTransactionManager 将是要走的路。

但以下代码将无法工作,因为 ChainedKafkaTransactionManager 需要 PlatformTransactionManager 类型的事务管理器。所以不接受参数r2dbcTransactionManager

    @Bean(name = ["chainedTransactionManager"])
    fun chainedTransactionManager(
        r2dbcTransactionManager: R2dbcTransactionManager,
        kafkaTransactionManager: KafkaTransactionManager<*, *>
    ) = ChainedKafkaTransactionManager(kafkaTransactionManager, r2dbcTransactionManager)

还有其他方法可以实现吗?

【问题讨论】:

    标签: spring-webflux project-reactor r2dbc reactive-kafka reactor-kafka


    【解决方案1】:

    为 Kafka 消费者链接事务是没有意义的。仅适用于发布者,即传出消息。

    但您应该确保不要多次处理同一消息。

    @EventListener(ApplicationStartedEvent::class)
    fun onMyEvent() {
      kafkaReceiver.receive()
        // Make sure to have unique index on (topic, partition, offset)
        // so you receive a ConstraintViolationException
        .flatMap { r ->
          val msg = ConsumedMessage(r.topic(), r.partition(), r.offset())
          consumedMessagesRepository.save(msg).thenReturn(r)
        }
        .onErrorContinue {ex, r -> log.warn("Duplicate msg") }
        .flatMap { r ->
          myService.deleteSomethingFromDbById(r.value().myId)
            .thenReturn(r)
        }
        .flatMap { r ->
          r.receiverOffset().commit()
        }
        .subscribe()
    }
    

    【讨论】:

      猜你喜欢
      • 2020-03-07
      • 1970-01-01
      • 2020-09-07
      • 1970-01-01
      • 2019-11-30
      • 2018-05-01
      • 2020-12-26
      • 2021-12-13
      • 1970-01-01
      相关资源
      最近更新 更多