kafka自定義分區(qū)器使用詳解
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)文章
阿里巴巴 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ù)的瀏覽記錄)(一)
這篇文章主要為大家詳細(xì)介紹了JAVAEE model1模型實(shí)現(xiàn)商品瀏覽記錄,去除重復(fù)的瀏覽記錄,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2016-11-11
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)前端傳參的操作,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2021-06-06

