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

kafka自定義分區(qū)器使用詳解

 更新時(shí)間:2025年11月19日 11:20:38   作者:princeAladdin  
本文介紹了如何根據(jù)企業(yè)需求自定義Kafka分區(qū)器,只需實(shí)現(xiàn)Partitioner接口并重寫(xiě)partition()方法,示例中,包含"cuihaida"的數(shù)據(jù)發(fā)送到0號(hào)分區(qū),否則發(fā)送到1號(hào)分區(qū),在生產(chǎn)者配置中添加分區(qū)器參數(shù)即可使用

kafka自定義分區(qū)器

根據(jù)企業(yè)需求,自己重新實(shí)現(xiàn)分區(qū)器

只需要定義類實(shí)現(xiàn)Partitioner接口,然后重寫(xiě)partition()方法即可

假設(shè)現(xiàn)在有一個(gè)需求

發(fā)送過(guò)來(lái)的數(shù)據(jù)中如果包含cuihaida,就發(fā)往0號(hào)分區(qū),不包含cuihaida,就發(fā)往1號(hào)分區(qū)

package com.example.kafkademo.producer;

import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;

import java.util.Map;

/**
 * 1. 實(shí)現(xiàn)接口Partitioner
 * 2. 實(shí)現(xiàn)3個(gè)方法:partition,close,configure
 * 3. 編寫(xiě)partition方法,返回分區(qū)號(hào)
 */
public class MyPartitioner implements Partitioner {

    /**
     * 重寫(xiě)這個(gè)方法
     * @param topic 主題
     * @param key 消息的key
     * @param keyBytes 消息的key序列化后的字節(jié)數(shù)組
     * @param value 消息的值
     * @param valueBytes 消息的值序列化后的字節(jié)數(shù)組
     * @param cluster 集群元數(shù)據(jù)可以查看分區(qū)信息
     * @return 信息對(duì)應(yīng)的分區(qū)
     */
    @Override
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
        // 獲取消息
        String msgValue = value.toString();
        // 發(fā)送過(guò)來(lái)的數(shù)據(jù)中如果包含cuihaida,就發(fā)往0號(hào)分區(qū),不包含cuihaida,就發(fā)往1號(hào)分區(qū)
        return msgValue.contains("cuihaida") ? 0 : 1;
    }

    @Override
    public void close() {

    }

    @Override
    public void configure(Map<String, ?> map) {

    }
}

使用分區(qū)器的方法

在生產(chǎn)者的配置中添加分區(qū)器參數(shù)

package com.example.kafkademo.util;

import org.apache.kafka.clients.producer.ProducerConfig;

import java.util.Properties;

public class CommonUtils {
    /**
     * kafka生產(chǎn)者配置配置
     * @return 配置內(nèi)容
     */
    public static Properties buildKafkaProperties() {
        // 1. 創(chuàng)建kafka生產(chǎn)者配置對(duì)象
        Properties properties = new Properties();
        // 2. 給kafka的配置對(duì)象添加信息
        properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "hadoop102:9092");
        // key, value初始化【必須有】
        properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
        properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
        
        // =========> 添加自定義分區(qū)器 <============
        properties.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, "com.example.kafkademo.producer.MyPartitioner")
        
        return properties;
    }
}

總結(jié)

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

相關(guān)文章

  • Java中的對(duì)稱加密詳解

    Java中的對(duì)稱加密詳解

    大家好,本篇文章主要講的是Java中的對(duì)稱加密詳解,感興趣的同學(xué)趕快來(lái)看一看吧,對(duì)你有幫助的話記得收藏一下
    2022-01-01
  • Netty粘包拆包及使用原理詳解

    Netty粘包拆包及使用原理詳解

    Netty是由JBOSS提供的一個(gè)java開(kāi)源框架,現(xiàn)為?Github上的獨(dú)立項(xiàng)目。Netty提供異步的、事件驅(qū)動(dòng)的網(wǎng)絡(luò)應(yīng)用程序框架和工具,用以快速開(kāi)發(fā)高性能、高可靠性的網(wǎng)絡(luò)服務(wù)器和客戶端程序,這篇文章主要介紹了Netty粘包拆包及使用原理
    2022-08-08
  • MyBatis Mapper代理使用方法詳解

    MyBatis Mapper代理使用方法詳解

    本文是小編日常收集整理的關(guān)于mybatis mapper代理使用方法知識(shí),通過(guò)本文還給大家提供有關(guān)開(kāi)發(fā)規(guī)范方面的知識(shí)點(diǎn),本文介紹的非常詳細(xì),具有參考借鑒價(jià)值,感興趣的朋友一起看下吧
    2016-08-08
  • Jetbrains系列產(chǎn)品重置試用思路詳解

    Jetbrains系列產(chǎn)品重置試用思路詳解

    這篇文章主要介紹了Jetbrains系列產(chǎn)品重置試用思路詳解,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2021-01-01
  • 阿里巴巴 Sentinel + InfluxDB + Chronograf 實(shí)現(xiàn)監(jiān)控大屏

    阿里巴巴 Sentinel + InfluxDB + Chronograf 實(shí)現(xiàn)監(jiān)控大屏

    這篇文章主要介紹了阿里巴巴 Sentinel + InfluxDB + Chronograf 實(shí)現(xiàn)監(jiān)控大屏,本文通過(guò)實(shí)例代碼給大家介紹的非常詳細(xì),具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2019-09-09
  • JAVAEE model1模型實(shí)現(xiàn)商品瀏覽記錄(去除重復(fù)的瀏覽記錄)(一)

    JAVAEE model1模型實(shí)現(xiàn)商品瀏覽記錄(去除重復(fù)的瀏覽記錄)(一)

    這篇文章主要為大家詳細(xì)介紹了JAVAEE model1模型實(shí)現(xiàn)商品瀏覽記錄,去除重復(fù)的瀏覽記錄,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2016-11-11
  • Java 二叉樹(shù)遍歷特別篇之Morris遍歷

    Java 二叉樹(shù)遍歷特別篇之Morris遍歷

    二叉樹(shù)的遍歷(traversing binary tree)是指從根結(jié)點(diǎn)出發(fā),按照某種次序依次訪問(wèn)二叉樹(shù)中所有的結(jié)點(diǎn),使得每個(gè)結(jié)點(diǎn)被訪問(wèn)依次且僅被訪問(wèn)一次。四種遍歷方式分別為:先序遍歷、中序遍歷、后序遍歷、層序遍歷
    2021-11-11
  • SpringBoot 啟動(dòng)報(bào)錯(cuò)Unable to connect to Redis server: 127.0.0.1/127.0.0.1:6379問(wèn)題的解決方案

    SpringBoot 啟動(dòng)報(bào)錯(cuò)Unable to connect to 

    這篇文章主要介紹了SpringBoot 啟動(dòng)報(bào)錯(cuò)Unable to connect to Redis server: 127.0.0.1/127.0.0.1:6379問(wèn)題的解決方案,文中通過(guò)圖文結(jié)合的方式給大家講解的非常詳細(xì),對(duì)大家解決問(wèn)題有一定的幫助,需要的朋友可以參考下
    2024-10-10
  • ConstraintValidator類如何實(shí)現(xiàn)自定義注解校驗(yàn)前端傳參

    ConstraintValidator類如何實(shí)現(xiàn)自定義注解校驗(yàn)前端傳參

    這篇文章主要介紹了ConstraintValidator類實(shí)現(xiàn)自定義注解校驗(yàn)前端傳參的操作,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-06-06
  • 一名Java高級(jí)工程師需要學(xué)什么?

    一名Java高級(jí)工程師需要學(xué)什么?

    作為一名Java高級(jí)工程師需要學(xué)什么?如何成為一名合格的工程師,這篇文章給了你較為詳細(xì)的答案,需要的朋友可以參考下
    2017-08-08

最新評(píng)論

望城县| 邻水| 汤阴县| 金阳县| 四川省| 石首市| 桐庐县| 曲阳县| 元氏县| 包头市| 东乌珠穆沁旗| 黄冈市| 淄博市| 周宁县| 清镇市| 齐齐哈尔市| 阿克| 海阳市| 岗巴县| 邓州市| 尖扎县| 赤峰市| 亳州市| 呼玛县| 巴中市| 黔南| 涞水县| 瑞安市| 鄂托克旗| 泾源县| 大名县| 元朗区| 孟连| 灌云县| 苍南县| 宁夏| 孝昌县| 芜湖县| 四平市| 清镇市| 尤溪县|