kafka add partitions function「建议收藏」

kafka add partitions function「建议收藏」代码功能在java代码中调用scala接口addPartitions.使用场景在kafka中如果需要定制kafka-topic的管理,那么其中一个功能很可能会用到:增加partition数量。但是在kafka-1.0.x之上的版本的AdminUtils中预留了相关的apiaddPartitions,具体功能的实现可以参考下面源码(scala):/***Addparti…

大家好,又见面了,我是你们的朋友全栈君。

代码功能

在 java 代码中调用 scala 接口 addPartitions.

使用场景

在kafka中如果需要定制kafka-topic的管理,那么其中一个功能很可能会用到:增加partition数量

但是在kafka-1.0.x之上的版本的AdminUtils中预留了相关的api addPartitions,具体功能的实现可以参考下面源码(scala):

/** * Add partitions to existing topic with optional replica assignment * * @param zkUtils Zookeeper utilities * @param topic Topic for adding partitions to * @param existingAssignment A map from partition id to its assigned replicas * @param allBrokers All brokers in the cluster * @param numPartitions Number of partitions to be set * @param replicaAssignment Manual replica assignment, or none * @param validateOnly If true, validate the parameters without actually adding the partitions * @return the updated replica assignment */
 @deprecated("This method is deprecated and will be replaced by kafka.zk.AdminZkClient.", "1.1.0")
 def addPartitions(zkUtils: ZkUtils,
                   topic: String,
                   existingAssignment: Map[Int, Seq[Int]],
                   allBrokers: Seq[BrokerMetadata],
                   numPartitions: Int = 1,
                   replicaAssignment: Option[Map[Int, Seq[Int]]] = None,
                   validateOnly: Boolean = false): Map[Int, Seq[Int]] = {
   val existingAssignmentPartition0 = existingAssignment.getOrElse(0,
     throw new AdminOperationException(
       s"Unexpected existing replica assignment for topic '$topic', partition id 0 is missing. " +
         s"Assignment: $existingAssignment"))

   val partitionsToAdd = numPartitions - existingAssignment.size
   if (partitionsToAdd <= 0)
     throw new InvalidPartitionsException(
       s"The number of partitions for a topic can only be increased. " +
         s"Topic $topic currently has ${existingAssignment.size} partitions, " +
         s"$numPartitions would not be an increase.")

   replicaAssignment.foreach { proposedReplicaAssignment =>
     validateReplicaAssignment(proposedReplicaAssignment, existingAssignmentPartition0,
       allBrokers.map(_.id).toSet)
   }

   val proposedAssignmentForNewPartitions = replicaAssignment.getOrElse {
     val startIndex = math.max(0, allBrokers.indexWhere(_.id >= existingAssignmentPartition0.head))
     AdminUtils.assignReplicasToBrokers(allBrokers, partitionsToAdd, existingAssignmentPartition0.size,
       startIndex, existingAssignment.size)
   }
   val proposedAssignment = existingAssignment ++ proposedAssignmentForNewPartitions
   if (!validateOnly) {
     info(s"Creating $partitionsToAdd partitions for '$topic' with the following replica assignment: " +
       s"$proposedAssignmentForNewPartitions.")
     // add the combined new list
     AdminUtils.createOrUpdateTopicPartitionAssignmentPathInZK(zkUtils, topic, proposedAssignment, update = true)
   }
   proposedAssignment

 }

复制代码

虽然 java 和 scala 之间的契合性很强,但是上述的参数列表在 java 中实现并调用该接口的复杂度还是比较高的,本人也是花费了一定的时间才调试通过的。

具体实现

下面部分是在 java 中调用该接口的实现:

public class addPartitions {
   public static boolean addPartitions(String zookeeperUri, String topic, int partitionNum) {
   	boolean succeed = true;
       ZkClient zkClient = createZkClient(zookeeperUri);
       ZkUtils zkUtils = ZkUtils.apply(zkClient, JaasUtils.isZkSecurityEnabled());

       //get existing assignment
       Seq<String> names = JavaConverters.asScalaBufferConverter(Arrays.asList(topic)).asScala();
       scala.collection.mutable.Map<String, scala.collection.Map<Object, Seq<Object>>> assignment = (scala.collection.mutable.Map<String, scala.collection.Map<Object, Seq<Object>>>)zkUtils.getPartitionAssignmentForTopics(names);
       Map<String, scala.collection.Map<Object, Seq<Object>>> partitionaAssigmentMap = JavaConverters.mutableMapAsJavaMapConverter(assignment).asJava();

       //get all brokers metadata in the cluster
       Seq<BrokerMetadata> allBrokerMeta = AdminUtils.getBrokerMetadatas(zkUtils, RackAwareMode.Enforced$.MODULE$, scala.Option.apply(null));

       try {
           AdminUtils.addPartitions(zkUtils, topic, partitionaAssigmentMap.get(topic), allBrokerMeta, partitionNum, Option.empty(), false);
       } catch (InvalidPartitionsException exception) {
           System.out.println("exception", exception);
           succeed = false;
       }
       return succeed;
   }
}
复制代码

此外,推荐 java 和 scala 之间格式转化的神器类 JavaConverters .

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 举报,一经查实,本站将立刻删除。

发布者:全栈程序员-用户IM,转载请注明出处:https://javaforall.cn/153334.html原文链接:https://javaforall.cn

【正版授权,激活自己账号】: Jetbrains全家桶Ide使用,1年售后保障,每天仅需1毛

【官方授权 正版激活】: 官方授权 正版激活 支持Jetbrains家族下所有IDE 使用个人JB账号...

(0)


相关推荐

  • 《开源安全运维平台-OSSIM最佳实践》已经上市

    《开源安全运维平台-OSSIM最佳实践》已经上市经多年潜心研究开源技术,历时三年创作的《开源安全运维平台OSSIM最佳实践》一书即将出版。该书用100多万字记录了作者10多年的OSSIM研究应用成果,重点展示了开源安全管理平台OSSIM在大型企业网运维管理中的实践。国内目前也有各式各样的运维系统,经过笔者对比分析得出这些工具无论在功能上、性能上还是在安全和稳定性易用性上都无法跟OSSIM系统想媲美,而且很多国内的开源安全运维项目在发布几年后就逐步淡出了舞台,而OSSIM持续发展了十多年。

    2022年10月25日
  • roundup与int的区别_notifyall()和notify()区别

    roundup与int的区别_notifyall()和notify()区别 isInterrupted()和interrputed()方法的区别isInterrupted方法是实例方法,interrupted方法是静态方法。Thread.currentThread().isInterrupted()Thread.interrupted()首先说明:wait(),notify(),notifyAll()这些方法由java.lang.Object类提供

    2022年10月28日
  • 网络攻防实验之缓冲区溢出攻击

    网络攻防实验之缓冲区溢出攻击这个实验是网络攻防课程实验中的一个,但是目前我还没有完全搞懂代码,以后有机会来补。也欢迎大佬指点一、实验目的和要求通过实验掌握缓冲区溢出的原理,通过使用缓冲区溢出攻击软件模拟入侵远程主机理解缓冲区溢出危害性,并理解防范和避免缓冲区溢出攻击的措施。二、实验原理和实验环境实验原理:缓冲区溢出(BufferOverflow)是目前非常普遍而且危…

  • C语言流水灯程序_51流水灯c语言程序

    C语言流水灯程序_51流水灯c语言程序0x01是数字,十六进制的数字。其结果等效于1。在数学上就是1,只不过在计算机上用2进制和十六进制较多,所以用十六进制表示。if(i&0x01)printf("奇数\n");elseprintf("偶数\n");system("pause");.0x01代表十六进制数也就是十进制数的01,&是把这些数转化为二进制数然后进行按位与运算info>>(…

  • vscode前端插件配置[通俗易懂]

    vscode前端插件配置[通俗易懂]vscode前端最佳插件配置

  • 贝叶斯优化python包_贝叶斯优化

    贝叶斯优化python包_贝叶斯优化万壑松风知客来,摇扇抚琴待留声1.文起本篇文章记录通过Python调用第三方库,从而调用使用了贝叶斯优化原理的Hyperopt方法来进行超参数的优化选择。具体贝叶斯优化原理与相关介绍将在下一次文章中做较为详细的描述,可以参考这里。Hyperopt是Python的几个贝叶斯优化库中的一个。它使用TreeParzenEstimator(TPE),其它Python库还包括了S…

    2022年10月22日

发表回复

您的电子邮箱地址不会被公开。

关注全栈程序员社区公众号