kafka自定义分区器使用详解
更新时间:2025年11月19日 11:20:38 作者:princeAladdin
本文介绍了如何根据企业需求自定义Kafka分区器,只需实现Partitioner接口并重写partition()方法,示例中,包含"cuihaida"的数据发送到0号分区,否则发送到1号分区,在生产者配置中添加分区器参数即可使用
kafka自定义分区器
根据企业需求,自己重新实现分区器
只需要定义类实现Partitioner接口,然后重写partition()方法即可
假设现在有一个需求
发送过来的数据中如果包含cuihaida,就发往0号分区,不包含cuihaida,就发往1号分区
package com.example.kafkademo.producer;
import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;
import java.util.Map;
/**
* 1. 实现接口Partitioner
* 2. 实现3个方法:partition,close,configure
* 3. 编写partition方法,返回分区号
*/
public class MyPartitioner implements Partitioner {
/**
* 重写这个方法
* @param topic 主题
* @param key 消息的key
* @param keyBytes 消息的key序列化后的字节数组
* @param value 消息的值
* @param valueBytes 消息的值序列化后的字节数组
* @param cluster 集群元数据可以查看分区信息
* @return 信息对应的分区
*/
@Override
public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
// 获取消息
String msgValue = value.toString();
// 发送过来的数据中如果包含cuihaida,就发往0号分区,不包含cuihaida,就发往1号分区
return msgValue.contains("cuihaida") ? 0 : 1;
}
@Override
public void close() {
}
@Override
public void configure(Map<String, ?> map) {
}
}
使用分区器的方法
在生产者的配置中添加分区器参数
package com.example.kafkademo.util;
import org.apache.kafka.clients.producer.ProducerConfig;
import java.util.Properties;
public class CommonUtils {
/**
* kafka生产者配置配置
* @return 配置内容
*/
public static Properties buildKafkaProperties() {
// 1. 创建kafka生产者配置对象
Properties properties = new Properties();
// 2. 给kafka的配置对象添加信息
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");
// =========> 添加自定义分区器 <============
properties.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, "com.example.kafkademo.producer.MyPartitioner")
return properties;
}
}
总结
以上为个人经验,希望能给大家一个参考,也希望大家多多支持脚本之家。
相关文章
阿里巴巴 Sentinel + InfluxDB + Chronograf 实现监控大屏
这篇文章主要介绍了阿里巴巴 Sentinel + InfluxDB + Chronograf 实现监控大屏,本文通过实例代码给大家介绍的非常详细,具有一定的参考借鉴价值,需要的朋友可以参考下2019-09-09
JAVAEE model1模型实现商品浏览记录(去除重复的浏览记录)(一)
这篇文章主要为大家详细介绍了JAVAEE model1模型实现商品浏览记录,去除重复的浏览记录,具有一定的参考价值,感兴趣的小伙伴们可以参考一下2016-11-11
SpringBoot 启动报错Unable to connect to 
这篇文章主要介绍了SpringBoot 启动报错Unable to connect to Redis server: 127.0.0.1/127.0.0.1:6379问题的解决方案,文中通过图文结合的方式给大家讲解的非常详细,对大家解决问题有一定的帮助,需要的朋友可以参考下2024-10-10
ConstraintValidator类如何实现自定义注解校验前端传参
这篇文章主要介绍了ConstraintValidator类实现自定义注解校验前端传参的操作,具有很好的参考价值,希望对大家有所帮助。如有错误或未考虑完全的地方,望不吝赐教2021-06-06


最新评论