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

怎樣給Kafka新增分區(qū)

 更新時間:2022年12月27日 15:19:16   作者:KK架構(gòu)  
這篇文章主要介紹了怎樣給Kafka新增分區(qū)問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教

給Kafka新增分區(qū)

數(shù)據(jù)量猛增的時候,需要給 kafka 的 topic 新增分區(qū),增大處理的數(shù)據(jù)量,可以通過以下步驟

1、修改 topic 的分區(qū)

kafka-topics --zookeeper hadoop004:2181 --alter --topic flink-test-04 --partitions 3

2、遷移數(shù)據(jù)

生成遷移計劃,手動新建一個 json 文件

{
"topics": [
{"topic": "flink-test-03"}
],
"version": 1
}

生成遷移計劃

kafka-reassign-partitions --zookeeper hadoop004:2181 --topics-to-move-json-file topic.json --broker-list “120,121,122” --generate

Current partition replica assignment:

{"version":1,"partitions":[{"topic":"flink-test-02","partition":5,"replicas":[120]},{"topic":"flink-test-02","partition":0,"replicas":[121]},{"topic":"flink-test-02","partition":2,"replicas":[120]},{"topic":"flink-test-02","partition":1,"replicas":[122]},{"topic":"flink-test-02","partition":4,"replicas":[122]},{"topic":"flink-test-02","partition":3,"replicas":[121]}]}

新建一個文件reassignment.json,保存上邊這些信息

3、遷移

kafka-reassign-partitions --zookeeper hadoop004:2181 --reassignment-json-file reassignment.json --execute

4、驗證

kafka-reassign-partitions --zookeeper hadoop004:2181 --reassignment-json-file reassignment.json --verify

Kafka分區(qū)原理機制

分區(qū)結(jié)構(gòu)

kafka的消息總共是三層結(jié)構(gòu)

Topic(第一層結(jié)構(gòu),表示一個主題)-> Partition(分區(qū),每個消息可以有多個分區(qū)) -> 消息實例(具體的消息文本等等,一個消息實例只可能在一個分區(qū)里面,不會出現(xiàn)在多個分區(qū)中)


在這里插入圖片描述

分區(qū)優(yōu)點

分區(qū)其實是一個負載均衡的思想。如此設(shè)計能使每一個分區(qū)獨自處理單獨的讀寫請求,提高吞吐量。

分區(qū)策略

  • 輪詢策略Round-robin(未指定key新版本默認策略)
  • 隨機策略Randomness(老版本默認策略)
  • 消息鍵排序策略Key-ordering(指定了key,則使用該策略)
  • 根據(jù)地理位置進行分區(qū)
  • 自定義分區(qū) 需要在生產(chǎn)者端實現(xiàn)org.apache.kafka.clients.producer.Partitioner接口,并配置一下實現(xiàn)類的全限定名

根據(jù)分區(qū)策略實現(xiàn)消息的順序消費

可以只設(shè)置一個分區(qū),這樣子消息都是放在一個partition,肯定是先進先出進行消費,然而這種場景無法利用kafka多分區(qū)的高吞吐量以及負載均衡的優(yōu)勢。

將需要順序消費的消息設(shè)置key,這個時候根據(jù)默認的分區(qū)策略,kafka會將所有的相同的key放在一個partition上面,這樣既可以使用kafka的partition又可以實現(xiàn)順序消費。

默認分區(qū)策略源碼

/**
 * The default partitioning strategy:
 * <ul>
 * <li>If a partition is specified in the record, use it
 * <li>If no partition is specified but a key is present choose a partition based on a hash of the key
 * <li>If no partition or key is present choose a partition in a round-robin fashion
 */
public class DefaultPartitioner implements Partitioner {
    private final ConcurrentMap<String, AtomicInteger> topicCounterMap = new ConcurrentHashMap<>();
    public void configure(Map<String, ?> configs) {}
    /**
     * Compute the partition for the given record.
     *
     * @param topic The topic name
     * @param key The key to partition on (or null if no key)
     * @param keyBytes serialized key to partition on (or null if no key)
     * @param value The value to partition on or null
     * @param valueBytes serialized value to partition on or null
     * @param cluster The current cluster metadata
     */
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
        List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
        int numPartitions = partitions.size();
        if (keyBytes == null) {
            int nextValue = nextValue(topic);
            List<PartitionInfo> availablePartitions = cluster.availablePartitionsForTopic(topic);
            if (availablePartitions.size() > 0) {
                int part = Utils.toPositive(nextValue) % availablePartitions.size();
                return availablePartitions.get(part).partition();
            } else {
                // no partitions are available, give a non-available partition
                return Utils.toPositive(nextValue) % numPartitions;
            }
        } else {
            // hash the keyBytes to choose a partition
            return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;
        }
    }
    private int nextValue(String topic) {
        AtomicInteger counter = topicCounterMap.get(topic);
        if (null == counter) {
            counter = new AtomicInteger(ThreadLocalRandom.current().nextInt());
            AtomicInteger currentCounter = topicCounterMap.putIfAbsent(topic, counter);
            if (currentCounter != null) {
                counter = currentCounter;
            }
        }
        return counter.getAndIncrement();
    }
    public void close() {}
}

從類注釋當(dāng)中已經(jīng)很明顯的看出來分區(qū)邏輯

3. 如果指定了分區(qū),則使用指定分區(qū)

4. 如果沒有指定分區(qū),但是有key,則使用hash過的key放置消息

5. 如果沒有指定分區(qū),也沒有key,則使用輪詢

總結(jié)

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

相關(guān)文章

  • 帶你快速搞定java IO

    帶你快速搞定java IO

    這篇文章主要介紹了Java IO流 文件傳輸基礎(chǔ)的相關(guān)資料,非常不錯,具有參考借鑒價值,需要的朋友可以參考下,希望能給你帶來幫助
    2021-07-07
  • Java 函數(shù)編程詳細介紹

    Java 函數(shù)編程詳細介紹

    這篇文章主要介紹了Java函數(shù)式編程,lambda表達式可以被認為是一個匿名函數(shù),可以在函數(shù)接口的上下文中使用。函數(shù)接口是只指定一個抽象方法的接口,下面來看文章的詳細內(nèi)容,需要的朋友可以參考下
    2021-11-11
  • mybatis中foreach報錯:_frch_item_0 not found的解決方法

    mybatis中foreach報錯:_frch_item_0 not found的解決方法

    這篇文章主要給大家介紹了mybatis中foreach報錯:_frch_item_0 not found的解決方法,文章通過示例代碼介紹了詳細的解決方法,對大家具有一定的參考學(xué)習(xí)價值,需要的朋友們下面來一起看看吧。
    2017-06-06
  • 聊聊@value注解和@ConfigurationProperties注解的使用

    聊聊@value注解和@ConfigurationProperties注解的使用

    這篇文章主要介紹了@value注解和@ConfigurationProperties注解的使用,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-09-09
  • Java實現(xiàn)迅雷地址轉(zhuǎn)成普通地址實例代碼

    Java實現(xiàn)迅雷地址轉(zhuǎn)成普通地址實例代碼

    本篇文章主要介紹了Java實現(xiàn)迅雷地址轉(zhuǎn)成普通地址實例代碼,非常具有實用價值,有興趣的可以了解一下。
    2017-03-03
  • Java中的顯示鎖ReentrantLock使用與原理詳解

    Java中的顯示鎖ReentrantLock使用與原理詳解

    這篇文章主要介紹了Java中的顯示鎖ReentrantLock使用與原理詳解,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-11-11
  • Java中的反射機制詳解

    Java中的反射機制詳解

    這篇文章主要介紹了JAVA 反射機制的相關(guān)知識,文中講解的非常細致,代碼幫助大家更好的理解學(xué)習(xí),感興趣的朋友可以了解下
    2021-09-09
  • java判斷Long類型的方法和實例代碼

    java判斷Long類型的方法和實例代碼

    在本篇文章里小編給大家整理的是關(guān)于java判斷Long類型的方法和實例代碼,對此有需要的朋友們跟著學(xué)習(xí)參考下。
    2020-02-02
  • java實現(xiàn)五子棋小游戲

    java實現(xiàn)五子棋小游戲

    這篇文章主要介紹了java實現(xiàn)五子棋小游戲的相關(guān)資料,十分簡單實用,推薦給大家,需要的朋友可以參考下
    2015-03-03
  • IDEA2020.1同步系統(tǒng)設(shè)置到GitHub的方法

    IDEA2020.1同步系統(tǒng)設(shè)置到GitHub的方法

    這篇文章主要介紹了IDEA2020.1同步系統(tǒng)設(shè)置到GitHub的方法,本文通過圖文并茂的形式給大家介紹的非常詳細,對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2020-05-05

最新評論

白山市| 阿鲁科尔沁旗| 海宁市| 富民县| 汝城县| 腾冲县| 蓬溪县| 洪雅县| 汶川县| 土默特右旗| 金平| 屯留县| 扶绥县| 平湖市| 嘉黎县| 凤阳县| 绥中县| 达州市| 罗甸县| 金昌市| 德令哈市| 西峡县| 陇川县| 鹤庆县| 宜宾县| 什邡市| 恭城| 乌审旗| 温宿县| 侯马市| 道真| 澳门| 泌阳县| 阿合奇县| 万全县| 曲水县| 甘南县| 阳高县| 鹿邑县| 樟树市| 岐山县|