最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

java發(fā)送kafka事務(wù)消息的實(shí)現(xiàn)方法

 更新時(shí)間:2022年07月15日 09:52:58   作者:逆風(fēng)飛翔的小叔  
本文主要介紹了java發(fā)送kafka事務(wù)消息的實(shí)現(xiàn)方法,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧

前言

事務(wù)對(duì)java開發(fā)的同學(xué)來說并不陌生,我們使用事務(wù)的目的在于避免產(chǎn)生重復(fù)數(shù)據(jù)或者說利用數(shù)據(jù)存儲(chǔ)中間件的事務(wù)特性確保數(shù)據(jù)的精準(zhǔn)性,比如大家熟悉的mysql,我們在程序開始時(shí),只需要在程序中添加上事務(wù)注解即可

kafka客戶端事務(wù),直接使用客戶端提供的相關(guān)的API即可,和jdbc事務(wù)的使用很類似,主要包含下面5個(gè)API

// 1 初始化事務(wù)
void initTransactions();


// 2 開啟事務(wù)
void beginTransaction() throws ProducerFencedException;


// 3 在事務(wù)內(nèi)提交已經(jīng)消費(fèi)的偏移量(主要用于消費(fèi)者)
void sendOffsetsToTransaction(Map<TopicPartition, OffsetAndMetadata> offsets,
 String consumerGroupId) throws ProducerFencedException;


// 4 提交事務(wù)
void commitTransaction() throws ProducerFencedException;


// 5 放棄事務(wù)(類似于回滾事務(wù)的操作)
void abortTransaction() throws ProducerFencedException;

下面結(jié)合實(shí)際的代碼以及效果演示進(jìn)行說明

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
 
import java.util.Properties;
 
public class ProducerTransaction {
 
    public static void main(String[] args) throws Exception {
 
        // 1. 創(chuàng)建 kafka 生產(chǎn)者的配置對(duì)象
        Properties properties = new Properties();
        // 2. 給 kafka 配置對(duì)象添加配置信息:bootstrap.servers
        properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "IP:9092");
 
        properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
        properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
 
        // 設(shè)置事務(wù) id(必須),事務(wù) id 任意起名
        properties.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "transaction_id_0");
 
        // 3. 創(chuàng)建 kafka 生產(chǎn)者對(duì)象
        KafkaProducer<String, String> kafkaProducer = new KafkaProducer<String, String>(properties);
 
        // 初始化事務(wù)
        kafkaProducer.initTransactions();
        // 開啟事務(wù)
        kafkaProducer.beginTransaction();
        System.out.println("開始發(fā)送消息");
        try {
            // 4. 調(diào)用 send 方法,發(fā)送消息
            for (int i = 0; i < 5; i++) {
                // 發(fā)送消息
                kafkaProducer.send(new ProducerRecord<>("zcy222", "hello kafka " + i));
            }
            //int i = 1 / 0;
            // 提交事務(wù)
            kafkaProducer.commitTransaction();
        } catch (Exception e) {
            System.out.println(e);
            // 終止事務(wù)
            kafkaProducer.abortTransaction();
        } finally {
            // 5. 關(guān)閉資源
            kafkaProducer.close();
        }
    }
 
}

運(yùn)行上面的代碼,正常是可以發(fā)送到指定的topic下

接下來,我們將上面的代碼中的 1/0 放開,再次運(yùn)行程序,可以看到,程序中拋異常了,但是消息并沒有發(fā)送到kafka的broker,說明事務(wù)的配置生效了

 到此這篇關(guān)于java發(fā)送kafka事務(wù)消息的實(shí)現(xiàn)方法的文章就介紹到這了,更多相關(guān)java發(fā)送kafka事務(wù)消息內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Spring 使用JavaConfig實(shí)現(xiàn)配置的方法步驟

    Spring 使用JavaConfig實(shí)現(xiàn)配置的方法步驟

    這篇文章主要介紹了Spring 使用JavaConfig實(shí)現(xiàn)配置的方法步驟,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-01-01
  • 詳解Java刪除Map中元素java.util.ConcurrentModificationException”異常解決

    詳解Java刪除Map中元素java.util.ConcurrentModificationException”異常解決

    這篇文章主要介紹了詳解Java刪除Map中元素java.util.ConcurrentModificationException”異常解,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2021-01-01
  • JSP 開發(fā)之 releaseSession的實(shí)例詳解

    JSP 開發(fā)之 releaseSession的實(shí)例詳解

    這篇文章主要介紹了JSP 開發(fā)之 releaseSession的實(shí)例詳解的相關(guān)資料,需要的朋友可以參考下
    2017-07-07
  • SpringBoot自定義starter方式

    SpringBoot自定義starter方式

    本文介紹了如何創(chuàng)建一個(gè)自定義的Spring Boot Starter,以實(shí)現(xiàn)日志功能,通過使用SPI機(jī)制,可以在不修改啟動(dòng)類的情況下,實(shí)現(xiàn)自動(dòng)配置和功能導(dǎo)入,同時(shí),還討論了如何在自定義Starter中編寫必要的配置文件和注解,以確保功能的正確實(shí)現(xiàn)和配置的智能提示
    2025-02-02
  • 關(guān)于application.yml數(shù)據(jù)庫配置方式

    關(guān)于application.yml數(shù)據(jù)庫配置方式

    這篇文章主要介紹了關(guān)于application.yml數(shù)據(jù)庫配置方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2024-08-08
  • IDEA插件Statistic統(tǒng)計(jì)代碼快速分辨爛項(xiàng)目

    IDEA插件Statistic統(tǒng)計(jì)代碼快速分辨爛項(xiàng)目

    這篇文章主要為大家介紹了使用IDEA插件Statistic來統(tǒng)計(jì)項(xiàng)目代碼,幫助大家快速識(shí)別出爛項(xiàng)目,有需要的朋友可以借鑒參考下,希望能夠有所幫助
    2022-01-01
  • Spring?Session(分布式Session共享)實(shí)現(xiàn)示例

    Spring?Session(分布式Session共享)實(shí)現(xiàn)示例

    這篇文章主要介紹了Spring?Session(分布式Session共享)實(shí)現(xiàn)示例,文章內(nèi)容詳細(xì),需要的朋友可以參考下
    2023-01-01
  • 基于JPA的Repository使用詳解

    基于JPA的Repository使用詳解

    這篇文章主要介紹了JPA的Repository使用詳解,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-11-11
  • Spring?IOC容器Bean注解創(chuàng)建對(duì)象組件掃描

    Spring?IOC容器Bean注解創(chuàng)建對(duì)象組件掃描

    這篇文章主要為大家介紹了Spring?IOC容器Bean注解創(chuàng)建對(duì)象組件掃描,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-05-05
  • Spring AOP里的靜態(tài)代理和動(dòng)態(tài)代理用法詳解

    Spring AOP里的靜態(tài)代理和動(dòng)態(tài)代理用法詳解

    這篇文章主要介紹了 Spring AOP里的靜態(tài)代理和動(dòng)態(tài)代理用法詳解,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-07-07

最新評(píng)論

佛山市| 德清县| 买车| 吕梁市| 奎屯市| 革吉县| 赤壁市| 靖远县| 宜章县| 博白县| 垣曲县| 五华县| 宁蒗| 修武县| 榆林市| 宁安市| 德惠市| 平塘县| 合山市| 兴义市| 云南省| 高碑店市| 谢通门县| 洛浦县| 吉水县| 龙川县| 清流县| 常山县| 东阳市| 石林| 慈利县| 宜昌市| 元氏县| 宜兴市| 大石桥市| 青田县| 师宗县| 普陀区| 寻甸| 禹城市| 河东区|