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

消息隊列-kafka消費異常問題

 更新時間:2021年07月02日 10:47:47   作者:禿頭披風(fēng)俠_  
這篇文章主要給大家介紹了關(guān)于kafka的相關(guān)資料,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧

概述

在kafka中,或者是說在任何消息隊列中都有個消費順序的問題。為了保證一個隊列順序消費,當(dāng)當(dāng)中一個消息消費異常時,必將影響后續(xù)隊列消息的消費,這樣業(yè)務(wù)豈不是卡住了。比如筆者舉個最簡單的例子:我發(fā)送1-100的消息,在我的處理邏輯當(dāng)中 msg%5==0我就進行 int i=1/0操作,這必將拋異常,一直阻塞在msg=5上,后面6-100無法消費。下面筆者給出解決方案。

重試一定次數(shù)(消息丟失)

@KafkaHandler
    @KafkaListener(topics = {"quickstart-events"},groupId = "test-consumer-group-2", concurrency = "1")
    public void test6(String msg){
              businessProcess(msg);
            }
           private void businessProcess(String msg){
        System.out.println("接收到消息:" + msg + "--" + System.currentTimeMillis() + "---" + Thread.currentThread().hashCode());
       if (Integer.valueOf(msg) % 5 == 0) {
            int i = 1 / 0;
        }
    }

說明:如果讀者使用的是java客戶端,也就是spring進行實現(xiàn),那么在不做任何處理的情況下,會自動重試10次,然后消息會被直接處理掉。也就是說如果你的業(yè)務(wù)允許消息丟失,那么你不需要額外的編碼處理

加入到死訊隊列(消息不丟失)

消費端代碼:

//1.啟用手動提交offset
//2.配置errorHandler,用來加入到死訊隊列
//3.不管業(yè)務(wù)處理是否處理異常還是正常都提交offset
@KafkaHandler
    @KafkaListener(topics = {"quickstart-events"},groupId = "test-consumer-group-2",
            errorHandler ="kafkaListenerErrorHandler", concurrency = "1")
    public void test6(String msg,Acknowledgment ack){
        try {
            businessProcess(msg);
        }finally {
            //手動提交
            ack.acknowledge();
        }
    }
//1.專門處理死訊隊列消息,都是topicName+.DLT的主題
//2.死訊隊列里,只有消費成功的才提交offset,否則等待bug修復(fù)完上線,繼續(xù)處理
    @KafkaHandler
    @KafkaListener(topics = {"quickstart-events.DLT"},groupId = "test-consumer-group-2", concurrency = "1")
    public void test7(String msg,Acknowledgment ack){
        try {
            businessProcess(msg);
            ack.acknowledge();
        }catch (Exception e){
            e.printStackTrace();
        }
    }
//業(yè)務(wù)代碼
    private void businessProcess(String msg){
        System.out.println("接收到消息:" + msg + "--" + System.currentTimeMillis() + "---" + Thread.currentThread().hashCode());
        if (Integer.valueOf(msg) % 5 == 0) {
            int i = 1 / 0;
        }
    }

異常處理器

//1.向容器注冊一個KafkaListenerErrorHandler類型的bean
//2.該bean就是當(dāng)處理消息異常的時候,將消息加入到.DLT主題中
@Component("kafkaListenerErrorHandler")
public class KafkaListenerErrorHandlerTest implements KafkaListenerErrorHandler {
   @Autowired
    private KafkaTemplate<String, Object> kafkaTemplate;
    private static final String TOPIC_DLT=".DLT";
    @Override
    public Object handleError(Message<?> message, ListenerExecutionFailedException exception) {
        System.out.println("消費失敗消息:"+message.toString());
        //獲取消息處理異常主題
        MessageHeaders headers = message.getHeaders();
        String topic=headers.get("kafka_receivedTopic")+TOPIC_DLT;
        //放入死訊隊列
        kafkaTemplate.send(topic,message.getPayload());
        return message;
    }
}

效果圖

image.png 

說明:以上基本上就是使用死訊隊列的方案,也許讀者會覺得這樣編碼復(fù)雜度很高,但其實不用擔(dān)心,其實上面這些代碼基本上是使用死訊隊列的模板代碼,在成熟一點的公司,一般會使用上述代碼進行簡單封裝,這里筆者給個思路,有興趣同學(xué)可以實現(xiàn)一下。我們其實可以使用aop思想,進行自定義一個@EnableDLT這樣的注解去實現(xiàn),這樣上面這個方案使用起來是不是就簡單優(yōu)雅了。之前筆者在開發(fā)過程中使用過亞馬遜的消息隊列服務(wù),也不過是這樣實現(xiàn)罷了。

總結(jié)

本篇文章就到這里了,希望可以給你帶來一些幫助,也希望您能夠多多關(guān)注腳本之家的更多內(nèi)容!

相關(guān)文章

  • SpringBoot之如何正確、安全的關(guān)閉服務(wù)

    SpringBoot之如何正確、安全的關(guān)閉服務(wù)

    這篇文章主要介紹了SpringBoot之如何正確、安全的關(guān)閉服務(wù)問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-03-03
  • Java的stream流多個字段排序的實現(xiàn)

    Java的stream流多個字段排序的實現(xiàn)

    本文主要介紹了Java的stream流多個字段排序的實現(xiàn),主要是兩種方法,第一種是固定多個字段排序和第二種動態(tài)字段進行排序,具有一定的參考價值,感興趣的可以了解一下
    2023-10-10
  • JAVA中Object的常用方法

    JAVA中Object的常用方法

    JAVA中Object是所有對象的頂級父類,存在于java.lang包中,這個包不需要我們手動導(dǎo)包,本文通過實例代碼介紹JAVA中Object的常用方法,感興趣的朋友一起看看吧
    2023-11-11
  • 基于spring mvc請求controller訪問方式

    基于spring mvc請求controller訪問方式

    這篇文章主要介紹了spring mvc請求controller訪問方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-09-09
  • 詳解如何給SpringBoot部署的jar包瘦身

    詳解如何給SpringBoot部署的jar包瘦身

    這篇文章主要介紹了如何給SpringBoot部署的jar包瘦身,如今迭代發(fā)布是常有的事情,每次都上傳一個如此龐大的文件,會浪費很多時間,接下來小編就以一個小項目為例,來演示如何給jar包瘦身,需要的朋友可以參考下
    2023-07-07
  • SpringBoot中RabbitMQ集群的搭建詳解

    SpringBoot中RabbitMQ集群的搭建詳解

    單個的?RabbitMQ?肯定無法實現(xiàn)高可用,要想高可用,還得上集群。這篇文章主要介紹了SpringBoot中RabbitMQ集群的兩種模式的搭建:普通集群搭建和鏡像集群搭建,需要的朋友可以參考一下
    2021-12-12
  • Java中Swing類實例講解

    Java中Swing類實例講解

    這篇文章主要介紹了Java中Swing類實例講解,文中用代碼實例講解的很清楚,有需要的同學(xué)可以研究下
    2021-02-02
  • Java?RabbitMQ消息隊列詳解常見問題

    Java?RabbitMQ消息隊列詳解常見問題

    消息隊列是最古老的中間件之一,從系統(tǒng)之間有通信需求開始,就自然產(chǎn)生了消息隊列。本文告訴什么是消息隊列,為什么需要消息隊列,常見的消息隊列有哪些,RabbitMQ的部署和使用
    2022-07-07
  • java微信紅包實現(xiàn)算法

    java微信紅包實現(xiàn)算法

    這篇文章主要為大家詳細介紹了java微信紅包實現(xiàn)算法,列出紅包的核心算法,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2017-02-02
  • Java嵌入式開發(fā)的優(yōu)勢及有點總結(jié)

    Java嵌入式開發(fā)的優(yōu)勢及有點總結(jié)

    在本篇內(nèi)容里小編給大家整理了關(guān)于Java嵌入式開發(fā)的優(yōu)勢及相關(guān)知識點內(nèi)容,有興趣的朋友們學(xué)習(xí)下。
    2022-11-11

最新評論

霸州市| 合川市| 噶尔县| 逊克县| 民勤县| 金秀| 太保市| 乌兰浩特市| 临泽县| 红原县| 盐边县| 当阳市| 五莲县| 建阳市| 奉化市| 敖汉旗| 思茅市| 阳春市| 余干县| 建德市| 密云县| 湟源县| 泰顺县| 曲靖市| 铜梁县| 通道| 五莲县| 香格里拉县| 威远县| 枞阳县| 长白| 垦利县| 鄯善县| 庐江县| 聂荣县| 积石山| 安陆市| 六枝特区| 伊通| 调兵山市| 桂东县|