欧美bbbwbbbw肥妇,免费乱码人妻系列日韩,一级黄片

流式圖表拒絕增刪改查之kafka核心消費(fèi)邏輯下篇

 更新時(shí)間:2023年04月12日 15:18:59   作者:在下uptown  
這篇文章主要為大家介紹了流式圖表拒絕增刪改查之kafka核心消費(fèi)邏輯講解的下篇,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

前篇回顧

kafka消費(fèi)者線(xiàn)程

突擊檢查八股文,實(shí)現(xiàn)線(xiàn)程的方法有哪些?嗯?沒(méi)復(fù)習(xí)是吧,行沒(méi)關(guān)系,那感謝參加本次面試哈。

常用的幾種方式分別是:

  • 繼承Thread類(lèi),重寫(xiě)run方法
  • 實(shí)現(xiàn)Runbale接口,重寫(xiě)run方法
  • 實(shí)現(xiàn)Callable接口,重寫(xiě)call方法

這里我們直接創(chuàng)捷出一個(gè)任務(wù)類(lèi)實(shí)現(xiàn)Runable方法,重寫(xiě)run方法,一個(gè)線(xiàn)程當(dāng)作一個(gè)kafka client,所以要在任務(wù)類(lèi)中聲明一個(gè)KafkaConsumer的成員變量,另外創(chuàng)建任務(wù)需要指定當(dāng)前任務(wù)的名稱(chēng)也就是線(xiàn)程名,還有要監(jiān)聽(tīng)的topic主題。

private KafkaConsumer<String, String> consumer;
private String topic;
private String threadName;

name和topic通過(guò)構(gòu)造方法傳進(jìn)來(lái),同時(shí)在構(gòu)造方法里完成對(duì)client的初始化操作。

/**
    * 封裝必要信息
    * @param bootServer 生產(chǎn)者ip
    * @param groupId 分組信息
    * @param topic  訂閱主題
    */
   public KafkaConsumerRunnable(String bootServer, String groupId, String topic) {
       this.topic = topic;
       Properties props = new Properties();
       props.put("bootstrap.servers", bootServer);
       props.put("group.id", groupId);
       props.put("enable.auto.commit", "false");
       props.put("auto.offset.reset", "latest");
       props.put("max.poll.records", 5);
       props.put("session.timeout.ms", "60000");
       props.put("max.poll.interval.ms", 300000);
       props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");  //鍵反序列化方式
       props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
       this.consumer = new KafkaConsumer&lt;&gt;(props);
   }

這里封裝kafka client的必要信息,入?yún)ootServer為kafka集群ip,groupId為threadName,我們規(guī)定一個(gè)線(xiàn)程為一個(gè)kafka消費(fèi)鏈接,消費(fèi)一個(gè)topic。

上一篇線(xiàn)程池保證了任務(wù)不會(huì)輕易掛掉,就算掛掉了也會(huì)重新提交,所以為了節(jié)省資源不做所謂的同groupId的負(fù)載操作。session.timeout.ms和max.poll.interval.ms可以根據(jù)當(dāng)前的kafka資源靈活配置,不然可能會(huì)引發(fā)一些reblance。

enable.auto.commit設(shè)置為false,手動(dòng)提交offset,auto.offset.reset這塊由于業(yè)務(wù)特殊,本來(lái)就是流式圖表瞬時(shí)的展示,如果真的出現(xiàn)了數(shù)據(jù)丟失那就丟了吧,從最新的數(shù)據(jù)讀取。

接下來(lái)只需要處理下消費(fèi)邏輯,consumer.subscribe(Collections.singletonList(this.topic))開(kāi)始訂閱監(jiān)聽(tīng)kafka數(shù)據(jù),搞一個(gè)while true不斷的消費(fèi)數(shù)據(jù),try catch只需要對(duì)WakeupException做處理,kafka客戶(hù)端會(huì)在關(guān)閉的時(shí)候拋出WakeupException異常。

finally里提交offset,無(wú)論這條offset對(duì)應(yīng)的數(shù)據(jù)消費(fèi)成功還是失敗都是消費(fèi)過(guò)了,失敗了就過(guò)去了。

   @Override
   public void run() {
   consumer.subscribe(Collections.singletonList(this.topic));
   String key = "stream_chart:" + this.name;
   Thread.currentThread().setName(key);
   try {
      while (true) {
         ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1));
         // 如果隊(duì)列中沒(méi)有消息 等待KAFKA_TIME_OUT后調(diào)用poll,如果有消息立即消費(fèi)
         for (ConsumerRecord<String, String> record : records) {
            String value = record.value();
            log.info("線(xiàn)程 {} 消費(fèi)kafka數(shù)據(jù) -> {} \n", Thread.currentThread().getName(), value);
            RedisConfig.getRedisTemplate().opsForZSet().add(key, value, Instant.now().getEpochSecond() * 1000);
         }
      }
   } catch (WakeupException e) {
      log.info("ignore for shutdown", e);
   } finally {
      consumer.commitAsync();
   }
}

我們消費(fèi)到數(shù)據(jù)直接放到redis的zset結(jié)構(gòu)里,當(dāng)前的時(shí)間戳作為score,最后留一個(gè)關(guān)閉客戶(hù)端的后門(mén)

// 退出后關(guān)掉客戶(hù)端
public void shutDown() {
   consumer.wakeup();
}

任務(wù)提交

任務(wù)提交這塊只需要在業(yè)務(wù)service中注入線(xiàn)程池,創(chuàng)建對(duì)應(yīng)的KafkaRunable任務(wù)封裝對(duì)應(yīng)的信息,執(zhí)行execute即可。

這里有個(gè)坑需要注意下,第二次突擊檢查八股文,線(xiàn)程池提交方法submitexecute的區(qū)別說(shuō)一下。不知道的立刻去熟讀并背誦。

public class TestTheadPool {
    public static void main(String[] args) {
        ExecutorService executorService= Executors.newFixedThreadPool(1);
        executorService.submit(new task("submit"));
        executorService.execute(new task("execute"));
    }
}
class task implements  Runnable{
    private String name;
    public task(String name) {
        this.name = name;
    }
    @Override
    public void run() {
        System.out.println(this.name + " start task");
        int i=1/0;
    }
}

熟悉的同學(xué)通過(guò)示例代碼可以看出來(lái),submit提交的線(xiàn)程不會(huì)拋出異常代碼,只有獲取Future返回值并執(zhí)行g(shù)et方法才會(huì)捕獲到異常。這塊涉及到異步的東西不再贅述

try {
    Future<?> submit = executorService.submit(new task("submit"));
    submit.get();
} catch (InterruptedException e) {
    e.printStackTrace();
} catch (ExecutionException e) {
    e.printStackTrace();
}

所以我們要使用execute執(zhí)行,不然kafka消費(fèi)線(xiàn)程里消費(fèi)失敗了攔截不到就不會(huì)被重新提交,導(dǎo)致線(xiàn)程掛掉。

以上就是流式圖表拒絕增刪改查之kafka核心消費(fèi)邏輯下篇的詳細(xì)內(nèi)容,更多關(guān)于kafka消費(fèi)流式圖表的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • 基于java實(shí)現(xiàn)租車(chē)管理系統(tǒng)

    基于java實(shí)現(xiàn)租車(chē)管理系統(tǒng)

    這篇文章主要為大家詳細(xì)介紹了基于java實(shí)現(xiàn)租車(chē)管理系統(tǒng),文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2020-12-12
  • Java實(shí)現(xiàn)二維碼、條形碼功能(案例代碼)

    Java實(shí)現(xiàn)二維碼、條形碼功能(案例代碼)

    ZXing是一個(gè)開(kāi)放源碼的,用Java實(shí)現(xiàn)的多種格式的1D/2D條碼圖像處理庫(kù),它包含了聯(lián)系到其他語(yǔ)言的端口,Zxing可以實(shí)現(xiàn)使用手機(jī)的內(nèi)置的攝像頭完成條形碼的掃描及解碼,這篇文章主要介紹了Java實(shí)現(xiàn)二維碼、條形碼等功能,需要的朋友可以參考下
    2024-01-01
  • SpringBoot項(xiàng)目啟動(dòng)后再請(qǐng)求遠(yuǎn)程接口的解決方式

    SpringBoot項(xiàng)目啟動(dòng)后再請(qǐng)求遠(yuǎn)程接口的解決方式

    Spring?Boot是由Pivotal團(tuán)隊(duì)提供的全新框架,其設(shè)計(jì)目的是用來(lái)簡(jiǎn)化Spring應(yīng)用的創(chuàng)建、運(yùn)行、調(diào)試、部署等,這篇文章主要介紹了SpringBoot項(xiàng)目啟動(dòng)后再請(qǐng)求遠(yuǎn)程接口的實(shí)現(xiàn)方式?,需要的朋友可以參考下
    2023-02-02
  • Java如何實(shí)現(xiàn)長(zhǎng)圖文生成的示例代碼

    Java如何實(shí)現(xiàn)長(zhǎng)圖文生成的示例代碼

    這篇文章主要介紹了Java如何實(shí)現(xiàn)長(zhǎng)圖文生成的示例代碼,小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2017-08-08
  • 淺談mybatis返回單一對(duì)象或?qū)ο罅斜淼膯?wèn)題

    淺談mybatis返回單一對(duì)象或?qū)ο罅斜淼膯?wèn)題

    這篇文章主要介紹了淺談mybatis返回單一對(duì)象或?qū)ο罅斜淼膯?wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-08-08
  • SpringMVC自定義攔截器登錄檢測(cè)功能的實(shí)現(xiàn)代碼

    SpringMVC自定義攔截器登錄檢測(cè)功能的實(shí)現(xiàn)代碼

    這篇文章主要介紹了SpringMVC自定義攔截器登錄檢測(cè)功能的實(shí)現(xiàn),本文通過(guò)實(shí)例代碼給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2021-08-08
  • java web實(shí)現(xiàn)簡(jiǎn)單留言板功能

    java web實(shí)現(xiàn)簡(jiǎn)單留言板功能

    這篇文章主要為大家詳細(xì)介紹了java web實(shí)現(xiàn)簡(jiǎn)單留言板功能,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2020-11-11
  • IDEA之啟動(dòng)參數(shù),配置文件默認(rèn)參數(shù)的操作

    IDEA之啟動(dòng)參數(shù),配置文件默認(rèn)參數(shù)的操作

    這篇文章主要介紹了IDEA之啟動(dòng)參數(shù),配置文件默認(rèn)參數(shù)的操作,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2021-01-01
  • Java中Spring的Security使用詳解

    Java中Spring的Security使用詳解

    這篇文章主要介紹了Java中Spring的Security使用詳解,在web應(yīng)用開(kāi)發(fā)中,安全無(wú)疑是十分重要的,選擇Spring Security來(lái)保護(hù)web應(yīng)用是一個(gè)非常好的選擇,需要的朋友可以參考下
    2023-07-07
  • Java中String、StringBuffer、StringBuilder的區(qū)別介紹

    Java中String、StringBuffer、StringBuilder的區(qū)別介紹

    這篇文章主要介紹了Java中String、StringBuffer、StringBuilder的區(qū)別介紹,本文講解了可變與不可變、是否多線(xiàn)程安全、gBuilder與StringBuffer共同點(diǎn)等內(nèi)容,需要的朋友可以參考下
    2015-06-06

最新評(píng)論