【问题标题】:Sample spring transaction with jpa and kafka?使用 jpa 和 kafka 进行 Spring 交易示例?
【发布时间】:2020-06-28 05:26:15
【问题描述】:

spring boot 从 2.1.11 版本升级到 2.2.5 后,kafka 客户端在提交 jpa 事务之前正在向代理生成消息。在不使用链式 kafka 事务管理器的情况下,事务工作正常。 2.1.x 和 2.2.x 之间的事务中是否有任何向后兼容性更改?有人可以提供任何跨越 JPA 和 Kafka 的工作事务管理器吗?

我只使用以下 JPA 事务管理器

  @Bean
  @Primary
  public PlatformTransactionManager transactionManager(final EntityManagerFactory emf) {
    final JpaTransactionManager txManager = new JpaTransactionManager();
    txManager.setEntityManagerFactory(emf);
    return txManager;
  }

我使用以下属性进行 kafka 事务:

spring.cloud.stream.kafka.default.consumer.configuration.isolation.level: read_committed spring.cloud.stream.kafka.binder.transaction.transaction-id-前缀:xyz-0 spring.kafka.producer.transaction-id-prefix: xyz-0

【问题讨论】:

    标签: spring-boot kafka-producer-api


    【解决方案1】:

    我们可以用 jpa 和 kafka 创建链式交易。我们可以使用 BinderFactory 来获取 kafka binder,并创建 binder kafka transaction manager。

    @Bean(name = "chainedTransactionManager")
      @Primary
      public PlatformTransactionManager chainedTransactionManager(JpaTransactionManager jpaTM,
          BinderFactory binders) {
        Binder<MessageChannel,?,?> binder = binders.getBinder("kafka", MessageChannel.class);
        if (binder instanceof KafkaMessageChannelBinder) {
          ProducerFactory<byte[], byte[]> pf =
              ((KafkaMessageChannelBinder) binder).getTransactionalProducerFactory();
          KafkaTransactionManager<byte[], byte[]> ktm = new KafkaTransactionManager<>(pf);
          ktm.setTransactionSynchronization(
              AbstractPlatformTransactionManager.SYNCHRONIZATION_ON_ACTUAL_TRANSACTION);
          return new ChainedKafkaTransactionManager<Object, Object>(jpaTM, ktm);
        } else {
          return jpaTM;
        }
      }```
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-04-07
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-05-31
      • 2013-06-08
      • 2019-10-08
      相关资源
      最近更新 更多