【发布时间】:2020-08-26 11:43:17
【问题描述】:
我正在尝试使用 KafkaTemplate 将消息从事务发布到 Kafka:
@Autowired
KafkaTemplate<GenericRecord, GenericRecord> kafkaTemplate;
@Transactional
@RabbitListener(queues = "queueName")
void input(final List<Message> messages) {
for (Message msg : messages) {
PublishRequest request = prepareRequest(msg);
kafkaTemplate.sendDefault(request.getKey(), reguest.getValue());
}
transactionalDatabaseInserts();
}
但是当我这样做时,我得到了这个异常:
原因:java.lang.IllegalStateException:没有事务在 过程;可能的解决方案:在 template.executeInTransaction() 操作的范围,开始一个 在调用模板方法之前使用@Transactional 进行事务, 在使用一个侦听器容器启动的事务中运行 记录
KafkaTemplate 的配置:
@EnableTransactionManagement
@Configuration
public class KafkaConfig{
@Bean
KafkaTransactionManager<GenericRecord, GenericRecord> kafkaTransactionManager(final ProducerFactory<GenericRecord, GenericRecord> producerFactory) {
return new KafkaTransactionManager<>(producerFactory);
}
@Bean
KafkaTemplate<GenericRecord, GenericRecord> kafkaTemplate(final ProducerFactory<GenericRecord, GenericRecord> producerFactory) {
return new KafkaTemplate<>(producerFactory);
}
}
在我的application.yaml 中,我包含了:
spring.kafka.producer.transaction-id-prefix: tx-
我希望我的方法适用于 @Transactional 而不是 kafkaTemplate.executeInTransaction()。为什么我会遇到这个异常?
【问题讨论】:
标签: java spring spring-boot apache-kafka spring-kafka