【问题标题】:How to use EntityManager/Hibernate Transaction with @KafkaHandler如何通过 @KafkaHandler 使用 EntityManager/Hibernate Transaction
【发布时间】:2022-01-24 14:32:00
【问题描述】:

假设我有一个 @KafkaListener 类,其中有一个 @KafkaHandler 方法,它处理任何收到的消息并执行一些数据库操作。

我想对此类中如何以及何时提交(或回滚)数据库更改(即手动管理数据库事务)进行细粒度控制。无论 DB 事务结果如何,都可以提交消费的消息偏移量。

这是我所拥有的简化版本:

@Service
@RequiredArgsConstructor
@KafkaListener(
    topics = "${kafka.topic.foo}",
    groupId = "${spring.kafka.consumer.group-id-foo}",
    containerFactory = "kafkaListenerContainerFactoryFoo")
public class FooMessageConsumer {

  // ...
  private final EntityManager entityManager;

  @KafkaHandler
  public void handleMessage(FooMessage msg) {
    // ...
    handleDBOperations(msg);
    // ...
  }

  void handleDBOperations(msg) {
    try {
      entityManager.getTransaction().begin();

      // ...

      entityManager.getTransaction().commit();
    } catch (Exception e) {
      log.error(e.getLocalizedMessage(), e);
      entityManager.getTransaction().rollback();
    }
  }

}

当收到消息并调用entityManager.getTransaction().begin(); 时,这会导致异常:

java.lang.IllegalStateException: Not allowed to create transaction on shared EntityManager - use Spring transactions or EJB CMT instead

为什么我不能在这里创建交易?

如果我删除 EntityManager 并将 @Transactional 注释添加到具有 DB 操作的方法(尽管这不是我想要的),那么它会导致另一个异常:

TransactionRequiredException Executing an update/delete query

它似乎完全忽略了注释。这是否与拥有自己的事务管理的 Kafka 消费者有关?

简而言之,我在这里做错了什么以及如何以@KafkaHandler 方法管理数据库事务?

感谢任何帮助。 提前致谢。

【问题讨论】:

    标签: java spring-boot hibernate apache-kafka spring-kafka


    【解决方案1】:

    尝试使用 Springs TransactionTemplate: https://docs.spring.io/spring-framework/docs/3.0.0.M4/reference/html/ch10s06.html

    如果您的用例很简单,Springs 声明式事务管理还应该让您实现您要求的行为: https://docs.spring.io/spring-framework/docs/3.0.0.M3/reference/html/ch11s05.html

    【讨论】:

    • 是的,使用TransactionTemplate 解决了它。但是,让我问一下,为什么我不能只使用EntityManager 事务甚至@Transactional 注释?
    • 存在一个关于为什么以这种方式询问实体经理不起作用的问题:stackoverflow.com/a/42901969; @Transactional 是 spring 通过 aop 代理在内部实现的,所以它只能在调用注入的 spring bean 的带注释的公共方法时工作(你尝试过吗?):docs.spring.io/spring-framework/docs/4.2.x/…
    • 我不知道实体管理器的范围,这很有帮助。但是,我仍然不明白为什么 @Transactional 在它已经是 Spring 管理对象的公共方法时不起作用。我还了解了 Kafka TM,这可能是我的情况下 DB TM 不起作用的原因(因此是 TransactionRequiredException),即使注释在那里。我将尝试明确指定 TM 并更新帖子以包含可能的解决方案,但谢谢! 12
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-12-26
    • 2012-05-24
    • 2012-12-29
    • 2011-11-20
    • 1970-01-01
    相关资源
    最近更新 更多