【发布时间】: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