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

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

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

前篇回顧

kafka消費者線程

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

常用的幾種方式分別是:

  • 繼承Thread類,重寫run方法
  • 實現(xiàn)Runbale接口,重寫run方法
  • 實現(xiàn)Callable接口,重寫call方法

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

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

name和topic通過構(gòu)造方法傳進來,同時在構(gòu)造方法里完成對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ī)定一個線程為一個kafka消費鏈接,消費一個topic。

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

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

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

finally里提交offset,無論這條offset對應的數(shù)據(jù)消費成功還是失敗都是消費過了,失敗了就過去了。

   @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));
         // 如果隊列中沒有消息 等待KAFKA_TIME_OUT后調(diào)用poll,如果有消息立即消費
         for (ConsumerRecord<String, String> record : records) {
            String value = record.value();
            log.info("線程 {} 消費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();
   }
}

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

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

任務提交

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

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

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;
    }
}

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

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

所以我們要使用execute執(zhí)行,不然kafka消費線程里消費失敗了攔截不到就不會被重新提交,導致線程掛掉。

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

相關(guān)文章

  • 一文帶你入門JDK8新特性——Lambda表達式

    一文帶你入門JDK8新特性——Lambda表達式

    這篇文章主要介紹了JDK8新特性——Lambda表達式的相關(guān)資料,幫助大家更好的理解和學習JAVA開發(fā),感興趣的朋友可以了解下
    2020-08-08
  • java線性表排序示例分享

    java線性表排序示例分享

    這篇文章主要介紹了java線性表排序示例,需要的朋友可以參考下
    2014-03-03
  • java?class?name實例深入精講

    java?class?name實例深入精講

    這篇文章主要為大家介紹了java?class?name實例深入精講,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2023-09-09
  • Java實現(xiàn)微信掃碼登入的實例代碼

    Java實現(xiàn)微信掃碼登入的實例代碼

    這篇文章主要介紹了java實現(xiàn)微信掃碼登入功能,本文通過實例代碼給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2020-06-06
  • Java實現(xiàn)簡單的萬年歷

    Java實現(xiàn)簡單的萬年歷

    這篇文章主要為大家詳細介紹了Java實現(xiàn)簡單的萬年歷,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2021-04-04
  • java簡單實現(xiàn)斗地主發(fā)牌功能

    java簡單實現(xiàn)斗地主發(fā)牌功能

    這篇文章主要為大家詳細介紹了java簡單實現(xiàn)斗地主發(fā)牌功能,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2021-06-06
  • idea項目實現(xiàn)移除和添加git

    idea項目實現(xiàn)移除和添加git

    本文指導讀者如何從官網(wǎng)下載并安裝Git,以及在IDEA中配置Git的詳細步驟,首先,用戶需訪問Git官方網(wǎng)站下載適合自己操作系統(tǒng)的Git版本并完成安裝,接著,在IDEA中通過設置找到git.exe文件以配置Gi
    2024-10-10
  • Scala小程序詳解及實例代碼

    Scala小程序詳解及實例代碼

    這篇文章主要介紹了Scala 第一個Scala小程序詳解的相關(guān)資料,需要的朋友可以參考下
    2017-01-01
  • Java中值傳遞的深度分析

    Java中值傳遞的深度分析

    這篇文章主要給大家介紹了關(guān)于Java中值傳遞的相關(guān)資料,文中通過示例代碼介紹的非常詳細,對大家學習或者使用java具有一定的參考學習價值,需要的朋友們下面來一起學習學習吧
    2019-04-04
  • Mybatis逆工程的使用

    Mybatis逆工程的使用

    最近在學Mybatis,類似Hibernate,Mybatis也有逆工程可以直接生成代碼(mapping,xml,pojo),方便快速開發(fā)。這篇文章給大家介紹Mybatis逆工程的使用相關(guān)知識,感興趣的朋友一起看下吧
    2016-06-06

最新評論

尼勒克县| 鸡东县| 丹凤县| 同心县| 江津市| 祁东县| 中山市| 通城县| 吴桥县| 牙克石市| 吴桥县| 长治市| 台前县| 尉氏县| 尉氏县| 台湾省| 城步| 临西县| 济宁市| 广饶县| 三江| 河北区| 准格尔旗| 德阳市| 桃园市| 酉阳| 富民县| 自治县| 阳信县| 乡城县| 石嘴山市| 白银市| 襄城县| 苍山县| 蒲江县| 新郑市| 南昌县| 通州市| 麦盖提县| 德惠市| 洛隆县|