springboot使用kafka事務(wù)的示例代碼
先看下下面這種情況,程序都出錯了,按理說消息也不應(yīng)該成功
@GetMapping("/send")
public void test9(String message) {
kafkaTemplate.send(topic, message);
throw new RuntimeException("fail");
}但是執(zhí)行結(jié)果是發(fā)生了異常并且消息發(fā)送成功了:

Kafka 同數(shù)據(jù)庫一樣支持事務(wù),當發(fā)生異常的時候可以進行回滾,確保消息監(jiān)聽器不會接收到一些錯誤的或者不需要的消息。
kafka事務(wù)屬性是指一系列的生產(chǎn)者生產(chǎn)消息和消費者提交偏移量的操作在一個事務(wù),或者說是是一個原子操作),同時成功或者失敗。使用事務(wù)也很簡單,需要先開啟事務(wù)支持,然后再使用。
如何開啟事務(wù)
如果使用默認配置只需要在yml添加spring.kafka.producer.transaction-id-prefix配置來開啟事務(wù),之前沒有使用默認的配置,自定義的kafkaTemplate,那么需要在ProducerFactory中設(shè)置事務(wù)Id前綴開啟事務(wù)并將KafkaTransactionManager注入到spring中,看下KafkaProducerConfig完整代碼:
@Configuration
@EnableKafka
public class KafkaProducerConfig {
@Value("${kafka.producer.servers}")
private String servers;
@Value("${kafka.producer.retries}")
private int retries;
public Map<String,Object> producerConfigs(){
Map<String,Object> props = new HashMap<>();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, servers);
props.put(ProducerConfig.RETRIES_CONFIG,retries);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
// 配置分區(qū)策略
props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG,"com.example.springbootkafka.config.CustomizePartitioner");
// 配置生產(chǎn)者攔截器
props.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG,"com.example.springbootkafka.interceptor.CustomProducerInterceptor");
// 配置攔截器消息處理類
SendMessageInterceptorUtil sendMessageInterceptorUtil = new SendMessageInterceptorUtil();
props.put("interceptorUtil",sendMessageInterceptorUtil);
return props;
}
@Bean
public ProducerFactory<String,String> producerFactory(){
DefaultKafkaProducerFactory producerFactory = new DefaultKafkaProducerFactory(producerConfigs());
//設(shè)置事務(wù)Id前綴 開啟事務(wù)
producerFactory.setTransactionIdPrefix("tx-");
return producerFactory;
}
@Bean
public KafkaTemplate<String,String> kafkaTemplate(){
return new KafkaTemplate<>(producerFactory());
}
@Bean
public KafkaTransactionManager<Integer, String> kafkaTransactionManager(ProducerFactory<String, String> producerFactory) {
return new KafkaTransactionManager(producerFactory);
}
} 配置開啟事務(wù)后,使用大體有兩種方式,先記錄下第一種使用事務(wù)方式:使用 executeInTransaction 方法
直接看下代碼:
@GetMapping("/send11")
public void test11(String message) {
kafkaTemplate.executeInTransaction(operations ->{
operations.send(topic,message);
throw new RuntimeException("fail");
});
}當然你可以這么寫:
@GetMapping("/send11")
public void test11(String message) {
kafkaTemplate.executeInTransaction(new KafkaOperations.OperationsCallback(){
@Override
public Object doInOperations(KafkaOperations operations) {
operations.send(topic,message);
throw new RuntimeException("fail");
}
});
}啟動項目,訪問http://localhost:8080/send10?message=test10 結(jié)果如下:

如上:消費者沒打印消息,說明消息沒發(fā)送成功,并且前面會報錯org.apache.kafka.common.KafkaException: Failing batch since transaction was aborted 的錯誤,說明事務(wù)生效了。
第一種使用事務(wù)方式:使用 @Transactional 注解方式 直接在方法上加上@Transactional注解即可,看下代碼:
@GetMapping("/send12")
@Transactional
public void test12(String message) {
kafkaTemplate.send(topic, message);
throw new RuntimeException("fail");
}如果開啟的事務(wù),則后續(xù)發(fā)送消息必須使用@Transactional注解或者使用kafkaTemplate.executeInTransaction() ,否則拋出異常,異常信息如下:
貼下完整的異常吧:java.lang.IllegalStateException: No transaction is in process; possible solutions: run the template operation within the scope of a template.executeInTransaction() operation, start a transaction with @Transactional before invoking the template method, run in a transaction started by a listener container when consuming a record
到此這篇關(guān)于springboot使用kafka事務(wù)的文章就介紹到這了,更多相關(guān)springboot kafka事務(wù)內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Java8 LocalDateTime極簡時間日期操作小結(jié)
這篇文章主要介紹了Java8-LocalDateTime極簡時間日期操作整理,通過實例代碼給大家介紹了java8 LocalDateTime 格式化問題,需要的朋友可以參考下2020-04-04
Java基礎(chǔ)教程之static五大應(yīng)用場景
這篇文章主要給大家介紹了關(guān)于Java基礎(chǔ)教程之static五大應(yīng)用場景的相關(guān)資料,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習或者工作具有一定的參考學(xué)習價值,需要的朋友們下面來一起學(xué)習學(xué)習吧2019-06-06
IDEA Error:java: 無效的源發(fā)行版: 17錯誤
本文主要介紹了IDEA Error:java: 無效的源發(fā)行版: 17錯誤,這個錯誤是因為您的IDEA編譯器不支持Java 17版本,您需要更新您的IDEA編譯器或者將您的Java版本降級到IDEA支持的版本,本文就來詳細的介紹一下2023-08-08
Spring Boot應(yīng)用監(jiān)控的實戰(zhàn)教程
Spring Boot 提供運行時的應(yīng)用監(jiān)控和管理功能,下面這篇文章主要給大家介紹了關(guān)于Spring Boot應(yīng)用監(jiān)控的相關(guān)資料,文中通過示例代碼介紹的非常詳細,需要的朋友可以參考借鑒,下面隨著小編來一起學(xué)習學(xué)習吧2018-05-05

