【问题标题】:Transaction Synchronization in Spring KafkaSpring Kafka 中的事务同步
【发布时间】:2018-05-01 11:01:05
【问题描述】:

我想将 kafka 事务与存储库事务同步:

@Transactional
public void syncTransaction(){
  myRepository.save(someObject)
  kafkaTemplate.send(someEvent)
}

由于合并 (https://github.com/spring-projects/spring-kafka/issues/373) 并且根据文档,这是可能的。不过,我在理解和实现该功能时遇到了问题。 查看https://docs.spring.io/spring-kafka/reference/html/#transaction-synchronization 中的示例,我必须创建一个 MessageListenerContainer 来监听我自己的事件。 我还需要使用 KafkaTemplate 发送我的事件吗? MessageListenerContainer 是否禁止发送给 broker?

如果我理解正确,kafkaTemplate 和 kafkaTransactionManager 必须使用相同的 producerFactory,我必须在其中启用 Transaction 设置 transactionIdPrefix。在我的示例中,我必须将 messageListenerContainer 的 TransactionManager 设置为 DataSourceTransactionManager。对吗?

从我的角度来看,我通过 kafkaTemplate 发送事件看起来很奇怪,监听我自己的事件并再次使用 kafkaTemplate 转发事件。

如果我能获得一个简单的 kafka 事务与存储库事务同步的示例和解释,我真的会帮助我。

【问题讨论】:

    标签: spring-transactions spring-kafka


    【解决方案1】:

    如果侦听器容器配置了KafkaTransactionManager,则容器将创建一个生产者,任何下游 kafka 模板都将使用该生产者,并且容器将为您将偏移量发送到事务。

    如果容器有其他事务管理器,则容器无法发送偏移量,因为它无权访问生产者(或模板)。

    另一种解决方案是使用 @Transactional(使用数据源 TM)注释您的方法,并使用 kafka TM 配置容器。

    这样,您的 DB tx 将在线程返回到容器之前提交,然后容器会将偏移量发送到 kafka 事务并提交。

    有关示例,请参阅 the framework test cases

    【讨论】:

    • "这样,您的 DB tx 将在线程返回到容器之前提交,然后容器会将偏移量发送到 kafka 事务并提交它。" => 但是 DB 提交和 Kafka 提交真的是一个原子操作吗?听起来像是两笔交易接二连三地执行,但不像是一笔交易。
    • 是的,但事务同步也是如此——它在Dr. Dave Syer's excellent Javaworld Artucle "Distributed transactions in Spring, with and without XA" 中被称为“Best Efforts 1PC 模式”。 Kafka 不支持 XA,您必须处理 DB tx 可能在 Kafka tx 回滚时提交的可能性。
    • 是否可以通过使用chainedTransaction 以原子方式提交DB 和Kafka?我有一个任务,它的状态保存在数据库中,同时,我需要将数据库记录 id 放入 kafka 主题并与 spring-kafka 消费者一起使用。这是将记录保存在数据库中并将其 id 放入 Kafka 然后在消费者中,根据放入 Kafka 的数据库记录 id 获取一些记录的好习惯吗?
    • 它不是原子的;这是“尽力而为 1 阶段提交”——请参阅我在上面评论中引用的文章。另请参阅docs.spring.io/spring-kafka/docs/2.6.2/reference/html/… - 最好不要在 cmets 中提出新问题,尤其是当答案是 3 岁时。事情会随着时间而改变。
    【解决方案2】:

    @Eike Behrends 有一个 db + kafka 事务,你可以使用ChainedTransactionManager 并这样定义它:

    @Bean
    public KafkaTransactionManager kafkaTransactionManager() {
        KafkaTransactionManager ktm = new KafkaTransactionManager(producerFactory());;
        ktm.setTransactionSynchronization(AbstractPlatformTransactionManager.SYNCHRONIZATION_ON_ACTUAL_TRANSACTION);
        return ktm;
    }
    
    
    @Bean
    @Primary
    public JpaTransactionManager transactionManager(EntityManagerFactory em) {
        return new JpaTransactionManager(em);
    }
    
    @Bean(name = "chainedTransactionManager")
    public ChainedTransactionManager chainedTransactionManager(JpaTransactionManager jpaTransactionManager,
                                                               KafkaTransactionManager kafkaTransactionManager) {
        return new ChainedTransactionManager(kafkaTransactionManager, jpaTransactionManager);
    }
    

    你需要注释你的事务性db+kafka方法@Transactional("chainedTransactionManager")

    (您可以在 spring-kafka 项目中看到问题:https://github.com/spring-projects/spring-kafka/issues/433

    你说:

    从我的角度来看,我通过以下方式发送事件看起来很奇怪 kafkaTemplate,监听我自己的事件并使用 再次使用 kafkaTemplate。

    你试过了吗?如果是这样,你能提供一个例子吗?

    【讨论】:

      【解决方案3】:

      为了实现您的目标,您应该使用不同的“最终一致”方法,例如 CDC(变更数据捕获)。 Kafka 写入和任何其他系统(例如数据库)之间没有原子事务——也就是 XA 事务。当您拥有分布式服务(有些称为微服务)时,这是一个完整的范例,在您的情况下,这些服务可能通过生产/消费到/从 Kafka 主题进行通信。

      【讨论】:

        【解决方案4】:

        TL;DR: 只需使用 upsert / merge。

        无意中看到这个老话题,这么多年了,大家还在纠结。

        只想分享最简单最原生的方法来处理像kafka这样的系统。

        人们来这里寻求答案的真正问题是分布式事务的旧方法。而且大多数人希望将非事务性(kafka 将其功能称为事务,但实际上它们是“特殊的”)kafka 与一些 ACID 数据库同步。

        如果您的服务在 idempotent 环境中运行 - 下游的所有内容也应该是 idempotent

        只要确保您对底层存储的操作是idempontent,最简单的方法是 upsert / merge(取决于存储)。

        附: CDC是一个东西,但它需要更多的人工成本,并且在大多数典型情况下是不必要的。

        更多: 如果您想深入了解为什么 kafka “事务”如此特别,这里有一些很好的起点(在 eos 中进行了解释):

        【讨论】:

          猜你喜欢
          • 2020-03-07
          • 2019-10-03
          • 1970-01-01
          • 2018-08-03
          • 2015-06-08
          • 2019-03-16
          • 1970-01-01
          • 1970-01-01
          • 2012-03-03
          相关资源
          最近更新 更多