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

解決kafka消息堆積及分區(qū)不均勻的問題

 更新時間:2021年09月11日 16:02:25   作者:筏鏡  
這篇文章主要介紹了解決kafka消息堆積及分區(qū)不均勻的問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教

kafka消息堆積及分區(qū)不均勻的解決

我在環(huán)境中發(fā)現(xiàn)代碼里面的kafka有所延遲,查看kafka消息發(fā)現(xiàn)堆積嚴重,經(jīng)過檢查發(fā)現(xiàn)是kafka消息分區(qū)不均勻造成的,消費速度過慢。這里由自己在虛擬機上演示相關問題,給大家提供相應問題的參考思路。

這篇文章有點遺憾并沒重現(xiàn)分區(qū)不均衡的樣例和Warning: Consumer group ‘testGroup1' is rebalancing. 這里僅將正確的方式展示,等后續(xù)重現(xiàn)了在進行補充。

主要有兩個要點:

  • 1、一個消費者組只消費一個topic.
  • 2、factory.setConcurrency(concurrency);這里設置監(jiān)聽并發(fā)數(shù)為 部署單元節(jié)點*concurrency=分區(qū)數(shù)量

1、先在kafka消息中創(chuàng)建

對應分區(qū)數(shù)目的topic(testTopic2,testTopic3)testTopic1由代碼創(chuàng)建

./kafka-topics.sh --create --zookeeper 192.168.25.128:2181 --replication-factor 1 --partitions 2 --topic testTopic2

2、添加配置文件application.properties

kafka.test.topic1=testTopic1
kafka.test.topic2=testTopic2
kafka.test.topic3=testTopic3
kafka.broker=192.168.25.128:9092
auto.commit.interval.time=60000
#kafka.test.group=customer-test
kafka.test.group1=testGroup1
kafka.test.group2=testGroup2
kafka.test.group3=testGroup3
kafka.offset=earliest
kafka.auto.commit=false

session.timeout.time=10000
kafka.concurrency=2

3、創(chuàng)建kafka工廠

package com.yin.customer.config;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.config.KafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.listener.AbstractMessageListenerContainer;
import org.springframework.kafka.listener.ConcurrentMessageListenerContainer;
import org.springframework.kafka.listener.ContainerProperties;
import org.springframework.stereotype.Component;
import java.util.HashMap;
import java.util.Map;

/**
 * @author yin
 * @Date 2019/11/24 15:54
 * @Method
 */
@Configuration
@Component
public class KafkaConfig {
    @Value("${kafka.broker}")
    private String broker;
    @Value("${kafka.auto.commit}")
    private String autoCommit;

   // @Value("${kafka.test.group}")
    //private String testGroup;

    @Value("${session.timeout.time}")
    private String sessionOutTime;

    @Value("${auto.commit.interval.time}")
    private String autoCommitTime;

    @Value("${kafka.offset}")
    private String offset;
    @Value("${kafka.concurrency}")
    private Integer concurrency;

   @Bean
    KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactory(){
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        //監(jiān)聽設置兩個個分區(qū)
        factory.setConcurrency(concurrency);
        //打開批量拉取數(shù)據(jù)
        factory.setBatchListener(true);
        //這里設置的是心跳時間也是拉的時間,也就說每間隔max.poll.interval.ms我們就調用一次poll,kafka默認是300s,心跳只能在poll的時候發(fā)出,如果連續(xù)兩次poll的時候超過
        //max.poll.interval.ms 值就會導致rebalance
        //心跳導致GroupCoordinator以為本地consumer節(jié)點掛掉了,引發(fā)了partition在consumerGroup里的rebalance。
        // 當rebalance后,之前該consumer擁有的分區(qū)和offset信息就失效了,同時導致不斷的報auto offset commit failed。
        factory.getContainerProperties().setPollTimeout(3000);
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
        return factory;
    }

    private ConsumerFactory<String,String> consumerFactory() {
        return new DefaultKafkaConsumerFactory<String, String>(consumerConfigs());
    }

   @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> propsMap = new HashMap<>();
        //kafka的地址
        propsMap.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, broker);
        //是否自動提交 Offset
        propsMap.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, autoCommit);
        // enable.auto.commit 設置成 false,那么 auto.commit.interval.ms 也就不被再考慮
        //默認5秒鐘,一個 Consumer 將會提交它的 Offset 給 Kafka
        propsMap.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG,  5000);

        //這個值必須設置在broker configuration中的group.min.session.timeout.ms 與 group.max.session.timeout.ms之間。
        //zookeeper.session.timeout.ms 默認值:6000
        //ZooKeeper的session的超時時間,如果在這段時間內沒有收到ZK的心跳,則會被認為該Kafka server掛掉了。
        // 如果把這個值設置得過低可能被誤認為掛掉,如果設置得過高,如果真的掛了,則需要很長時間才能被server得知。
        propsMap.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, sessionOutTime);
        propsMap.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        propsMap.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        //組與組間的消費者是沒有關系的。
        //topic中已有分組消費數(shù)據(jù),新建其他分組ID的消費者時,之前分組提交的offset對新建的分組消費不起作用。
        //propsMap.put(ConsumerConfig.GROUP_ID_CONFIG, testGroup);

        //當創(chuàng)建一個新分組的消費者時,auto.offset.reset值為latest時,
        // 表示消費新的數(shù)據(jù)(從consumer創(chuàng)建開始,后生產(chǎn)的數(shù)據(jù)),之前產(chǎn)生的數(shù)據(jù)不消費。
        // https://blog.csdn.net/u012129558/article/details/80427016

        //earliest 當分區(qū)下有已提交的offset時,從提交的offset開始消費;無提交的offset時,從頭開始消費。
       // latest 當分區(qū)下有已提交的offset時,從提交的offset開始消費;無提交的offset時,消費新產(chǎn)生的該分區(qū)下的數(shù)據(jù)。

        propsMap.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, offset);
        //不是指每次都拉50條數(shù)據(jù),而是一次最多拉50條數(shù)據(jù)()
        propsMap.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 5);
        return propsMap;
    }
}

4、展示kafka消費者

@Component
public class KafkaConsumer {
    private static final Logger logger = LoggerFactory.getLogger(KafkaConsumer.class);

    @KafkaListener(topics = "${kafka.test.topic1}",groupId = "${kafka.test.group1}",containerFactory = "kafkaListenerContainerFactory")
    public void listenPartition1(List<ConsumerRecord<?, ?>> records,Acknowledgment ack) {
        logger.info("testTopic1 recevice a message size :{}" , records.size());

        try {
            for (ConsumerRecord<?, ?> record : records) {
                Optional<?> kafkaMessage = Optional.ofNullable(record.value());
                logger.info("received:{} " , record);
                if (kafkaMessage.isPresent()) {
                    Object message = record.value();
                    String topic = record.topic();
                    Thread.sleep(300);
                    logger.info("p1 topic is:{} received message={}",topic, message);
                }
            }
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            ack.acknowledge();
        }
    }

    @KafkaListener(topics = "${kafka.test.topic2}",groupId = "${kafka.test.group2}",containerFactory = "kafkaListenerContainerFactory")
    public void listenPartition2(List<ConsumerRecord<?, ?>> records,Acknowledgment ack) {
        logger.info("testTopic2 recevice a message size :{}" , records.size());

        try {
            for (ConsumerRecord<?, ?> record : records) {
                Optional<?> kafkaMessage = Optional.ofNullable(record.value());
                logger.info("received:{} " , record);
                if (kafkaMessage.isPresent()) {
                    Object message = record.value();
                    String topic = record.topic();
                    Thread.sleep(300);
                    logger.info("p2 topic :{},received message={}",topic, message);
                }
            }
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            ack.acknowledge();
        }
    }

    @KafkaListener(topics = "${kafka.test.topic3}",groupId = "${kafka.test.group3}",containerFactory = "kafkaListenerContainerFactory")
    public void listenPartition3(List<ConsumerRecord<?, ?>> records, Acknowledgment ack) {
        logger.info("testTopic3 recevice a message size :{}" , records.size());

        try {
            for (ConsumerRecord<?, ?> record : records) {
                Optional<?> kafkaMessage = Optional.ofNullable(record.value());
                logger.info("received:{} " , record);
                if (kafkaMessage.isPresent()) {
                    Object message = record.value();
                    String topic = record.topic();
                    logger.info("p3 topic :{},received message={}",topic, message);
                    Thread.sleep(300);
                }
            }
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            ack.acknowledge();
        }
    }
}

查看分區(qū)消費情況:

在這里插入圖片描述

kafka出現(xiàn)若干分區(qū)不消費的現(xiàn)象

近日,有用戶反饋kafka有topic出現(xiàn)某個消費組消費的時候,有幾個分區(qū)一直不消費消息,消息一直積壓(圖1)。除了一直積壓外,還有一個現(xiàn)象就是消費組一直在重均衡,大約每5分鐘就會重均衡一次。具體表現(xiàn)為消費分區(qū)的owner一直在改變(圖2)。

(圖1)

(圖2)

定位過程

業(yè)務側沒有報錯,同時kafka服務端日志也一切正常,同事先將消費組的機器滾動重啟,仍然還是那幾個分區(qū)沒有消費,之后將這幾個不消費的分區(qū)遷移至別的broker上,依然沒有消費。

還有一個奇怪的地方,就是每次重均衡后,不消費的那幾個分區(qū)的消費owner所在機器的網(wǎng)絡都有流量變化。按理說不消費應該就是拉取不到分區(qū)不會有流量的。于是讓運維去拉了下不消費的consumer的jstack日志。一看果然發(fā)現(xiàn)了問題所在。

從堆??矗琧onsumer已經(jīng)拉取到消息,然后就一直卡在處理消息的業(yè)務邏輯上。這說明kafka是沒有問題的,用戶的業(yè)務邏輯有問題。

consumer在拉取完一批消息后,就一直在處理這批消息,但是這批消息中有若干條消息無法處理,而業(yè)務又沒有超時操作或者異常處理導致進程一直處于消費中,無法去poll下一批數(shù)據(jù)。

又由于業(yè)務采用的是autocommit的offset提交方式,而根據(jù)源碼可知,consumer只有在下一次poll中才會自動提交上次poll的offset,所以業(yè)務一直在拉取同一批消息而無法更新offset。反映的現(xiàn)象就是該consumer對應的分區(qū)的offset一直沒有變,所以有積壓的現(xiàn)象。

至于為什么會一直在重均衡消費組的原因也很明了了,就是因為有消費者一直卡在處理消息的業(yè)務邏輯上,超過了max.poll.interval.ms(默認5min),消費組就會將該消費者踢出消費組,從而發(fā)生重均衡。

驗證

讓業(yè)務方去查證業(yè)務日志,驗證了積壓的這幾個分區(qū),總是在循環(huán)的拉取同一批消息。

解決方法

臨時解決方法就是跳過有問題的消息,將offset重置到有問題的消息之后。本質上還是要業(yè)務側修改業(yè)務邏輯,增加超時或者異常處理機制,最好不要采用自動提交offset的方式,可以手動管理。

以上為個人經(jīng)驗,希望能給大家一個參考,也希望大家多多支持腳本之家。

相關文章

  • JAVA的Dubbo如何實現(xiàn)各種限流算法

    JAVA的Dubbo如何實現(xiàn)各種限流算法

    Dubbo是一種高性能的Java RPC框架,廣泛應用于分布式服務架構中,在Dubbo中實現(xiàn)限流可以幫助服務在高并發(fā)場景下保持穩(wěn)定性和可靠性,常見的限流算法包括固定窗口算法、滑動窗口算法、令牌桶算法和漏桶算法,在Dubbo中集成限流器可以通過實現(xiàn)自定義過濾器來實現(xiàn)
    2025-01-01
  • Java實現(xiàn)Excel表單控件的添加與刪除

    Java實現(xiàn)Excel表單控件的添加與刪除

    本文通過Java代碼示例介紹如何在Excel表格中添加表單控件,包括文本框、單選按鈕、復選框、組合框、微調按鈕等,以及如何刪除Excel中的指定表單控件,需要的可以參考一下
    2022-05-05
  • MyBatis?Generator生成的$?sql是否存在注入風險詳解

    MyBatis?Generator生成的$?sql是否存在注入風險詳解

    這篇文章主要介紹了MyBatis?Generator生成的$?sql是否存在注入風險詳解,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-12-12
  • Java堆空間爆滿導致宕機的問題分析及解決

    Java堆空間爆滿導致宕機的問題分析及解決

    團隊有一個服務,一直運行的好好的,突然訪問異常了,先是請求超時,然后直接無法訪問,本文將給大家介紹Java堆空間爆滿導致宕機的問題分析及解決,需要的朋友可以參考下
    2024-02-02
  • SpringBoot的jar包如何啟動的實現(xiàn)

    SpringBoot的jar包如何啟動的實現(xiàn)

    本文主要介紹了SpringBoot的jar包如何啟動的實現(xiàn),文中根據(jù)實例編碼詳細介紹的十分詳盡,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2022-03-03
  • java微信server錄音下載到自己server

    java微信server錄音下載到自己server

    這篇文章主要為大家詳細介紹了java微信server錄音下載到自己server的相關代碼,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2017-05-05
  • springboot+zookeeper實現(xiàn)分布式鎖的示例代碼

    springboot+zookeeper實現(xiàn)分布式鎖的示例代碼

    本文主要介紹了springboot+zookeeper實現(xiàn)分布式鎖的示例代碼,文中根據(jù)實例編碼詳細介紹的十分詳盡,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2022-03-03
  • Java ForkJoin框架的原理及用法

    Java ForkJoin框架的原理及用法

    這篇文章主要介紹了Java ForkJoin框架的原理及用法,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友可以參考下
    2019-10-10
  • SpringBoot實現(xiàn)國際化的教程

    SpringBoot實現(xiàn)國際化的教程

    這篇文章主要介紹了SpringBoot實現(xiàn)國際化的教程,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2025-03-03
  • java使用freemarker模板生成html再轉為pdf

    java使用freemarker模板生成html再轉為pdf

    這篇文章主要為大家詳細介紹了java如何使用freemarker模板生成html,再利用iText將生成的HTML轉換為PDF文件,感興趣的小伙伴可以參考下
    2025-04-04

最新評論

福安市| 都江堰市| 报价| 平顺县| 仪陇县| 汝城县| 香河县| 永胜县| 大渡口区| 临沭县| 天门市| 鄂托克旗| 勃利县| 青冈县| 武平县| 左云县| 额尔古纳市| 奇台县| 永登县| 科尔| 沈丘县| 蓝田县| 凤庆县| 隆化县| 徐汇区| 库尔勒市| 赣榆县| 邢台市| 噶尔县| 理塘县| 卢湾区| 三门县| 台南市| 依安县| 永城市| 余姚市| 巴中市| 灵丘县| 金塔县| 呼玛县| 库伦旗|