当前位置:   article > 正文

springboot集成整合kafka-自定义分区策略、将消息发送到指定的分区partition_springboot集成kafka往指定分区中发送消息

springboot集成kafka往指定分区中发送消息

写在前面:各位看到此博客的小伙伴,如有不对的地方请及时通过私信我或者评论此博客的方式指出,以免误人子弟。多谢!

  kafka中每个topic可定义多个分区,那么生产者将消息发送到topic时,具体追加到哪个分区呢?上一篇我们已经通过这种方式:@KafkaListener(topics = {"mytopic"},topicPartitions = {@TopicPartition(topic = "mytopic",partitions = {"1","3"})})实现了只消费指定分区的消息,但并没有实现将消息发送到指定分区,这篇就记录一下kafka自定义分区策略-将消息发送到我们指定的partition。

其实不用自定义分区策略一样可以将消息发送到指定分区,看下send()方法:

可以看到有两个重载的send方法可以直接指定分区发送消息,这里就不测试了,直接看自定义分区策略。

在自定义分区策略前,先看下kafka默认是怎么的,访问http://localhost:8080/send1?message=test1,看下打印结果:

可以看到消息随机分布到topic的各个分区上,这是kafka默认的分区策略,轮询选出一个patition。

在自定义分区策略前,我们还应该知道消息的路由机制:

  1. 若发送消息时指定了分区(即自定义分区策略),则直接将消息append到指定分区;
  2. 若发送消息时未指定patition,但指定了key(kafka允许为每条消息设置一个key),则对key值进行hash计算,根据计算结果路由到指定分区,这种情况下可以保证同一个Key的所有消息都进入到相同的分区;
  3. patition 和 key 都未指定,则使用kafka默认的分区策略,轮询选出一个 patition;

测试一下为消息指定key时是否会将消息发送到指定分区,通过send的另一个方法send(String var1, K var2, V var3)测试一下:

代码如下:

  1. @GetMapping("/send3")
  2. public String send3(String message) {
  3. String key = "zxcvb";
  4. kafkaTemplate.send(topic,key, message).addCallback(
  5. success ->{
  6. String topic = success.getRecordMetadata().topic();
  7. int partition = success.getRecordMetadata().partition();
  8. long offset = success.getRecordMetadata().offset();
  9. System.out.println("topic:" + topic + " partition:" + partition + " offset:" + offset);
  10. },
  11. failure ->{
  12. String message1 = failure.getMessage();
  13. System.out.println(message1);
  14. }
  15. );
  16. return "success";
  17. }

访问http://localhost:8080/send3?message=test3  多次,打印结果如下:

可以看到消息都发送到partition1上了。

下面说下本篇的正题,自定义分区策略,自定义分区策略其实很简单,只要实现Partitioner接口,重写它的方法就好了,看下代码:

  1. @Component
  2. public class CustomizePartitioner implements Partitioner {
  3. /**
  4. * 自定义分区策略
  5. */
  6. @Override
  7. public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
  8. // 获取topic的分区列表
  9. List<PartitionInfo> partitionInfoList = cluster.availablePartitionsForTopic(topic);
  10. int partitionCount = partitionInfoList.size();
  11. int auditPartition = partitionCount - 1;
  12. return auditPartition;
  13. }
  14. /**
  15. * 在分区程序关闭时调用
  16. */
  17. @Override
  18. public void close() {
  19. System.out.println("colse ...");
  20. }
  21. /**
  22. * 做必要的初始化工作
  23. */
  24. @Override
  25. public void configure(Map<String, ?> configs) {
  26. System.out.println("init ...");
  27. }
  28. }

定义完分区策略,需要配置下分区策略,在KafkaProducerConfig中producerConfigs方法里追加如下一行代码即可:

props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG,"com.example.springbootkafka.config.CustomizePartitioner");

新增测试方法:

  1. @GetMapping("/send5")
  2. public String send5(String message) {
  3. kafkaTemplate.send(topic, message).addCallback(
  4. success ->{
  5. String topic = success.getRecordMetadata().topic();
  6. int partition = success.getRecordMetadata().partition();
  7. long offset = success.getRecordMetadata().offset();
  8. System.out.println("topic:" + topic + " partition:" + partition + " offset:" + offset);
  9. },
  10. failure ->{
  11. String message1 = failure.getMessage();
  12. System.out.println(message1);
  13. }
  14. );
  15. return "success";
  16. }

访问http://localhost:8080/send5?message=test5 多次,查看打印结果如下:

由打印结果来看,消息都发到partition4分区了,另外mytopic有5个分区,返回auditPartition为4,也没错。

声明:本文内容由网友自发贡献,转载请注明出处:【wpsshop博客】
推荐阅读
相关标签
  

闽ICP备14008679号