From e71e5a442d7bfd510f5a29afb206d702bc2b3880 Mon Sep 17 00:00:00 2001 From: wenbingshen Date: Fri, 12 Mar 2021 01:31:29 +0800 Subject: [PATCH 01/10] KAFKA-12454:Add ERROR logging on kafka-log-dirs when given brokerIds do not exist in current kafka cluster --- core/src/main/scala/kafka/admin/LogDirsCommand.scala | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/admin/LogDirsCommand.scala b/core/src/main/scala/kafka/admin/LogDirsCommand.scala index ad0307f31b9cc..affe518b33ca4 100644 --- a/core/src/main/scala/kafka/admin/LogDirsCommand.scala +++ b/core/src/main/scala/kafka/admin/LogDirsCommand.scala @@ -40,11 +40,16 @@ object LogDirsCommand { val opts = new LogDirsCommandOptions(args) val adminClient = createAdminClient(opts) val topicList = opts.options.valueOf(opts.topicListOpt).split(",").filter(!_.isEmpty) + val clusterBrokers: Array[Int] = adminClient.describeCluster().nodes().get().asScala.map(_.id()).toArray val brokerList = Option(opts.options.valueOf(opts.brokerListOpt)) match { case Some(brokerListStr) => brokerListStr.split(',').filter(!_.isEmpty).map(_.toInt) - case None => adminClient.describeCluster().nodes().get().asScala.map(_.id()).toArray + case None => clusterBrokers } + val nonExistBrokers: Array[Int] = brokerList.filterNot(brokerId => clusterBrokers.contains(brokerId)) + if (!nonExistBrokers.isEmpty) + System.err.println(s"The given node(s) does not exist from broker-list ${nonExistBrokers.mkString(",")}") + out.println("Querying brokers for log directories information") val describeLogDirsResult: DescribeLogDirsResult = adminClient.describeLogDirs(brokerList.map(Integer.valueOf).toSeq.asJava) val logDirInfosByBroker = describeLogDirsResult.allDescriptions.get().asScala.map { case (k, v) => k -> v.asScala } From 91fd3057a12d7b80f7698d49616efdb03477beb2 Mon Sep 17 00:00:00 2001 From: wenbingshen Date: Fri, 12 Mar 2021 01:36:34 +0800 Subject: [PATCH 02/10] KAFKA-12454:Add ERROR logging on kafka-log-dirs when given brokerIds do not exist in current kafka cluster --- core/src/main/scala/kafka/admin/LogDirsCommand.scala | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/admin/LogDirsCommand.scala b/core/src/main/scala/kafka/admin/LogDirsCommand.scala index affe518b33ca4..46193ef3248f7 100644 --- a/core/src/main/scala/kafka/admin/LogDirsCommand.scala +++ b/core/src/main/scala/kafka/admin/LogDirsCommand.scala @@ -47,8 +47,10 @@ object LogDirsCommand { } val nonExistBrokers: Array[Int] = brokerList.filterNot(brokerId => clusterBrokers.contains(brokerId)) - if (!nonExistBrokers.isEmpty) - System.err.println(s"The given node(s) does not exist from broker-list ${nonExistBrokers.mkString(",")}") + if (!nonExistBrokers.isEmpty) { + System.err.println(s"The given node(s) does not exist from broker-list ${nonExistBrokers.mkString(",")}") + sys.exit(1) + } out.println("Querying brokers for log directories information") val describeLogDirsResult: DescribeLogDirsResult = adminClient.describeLogDirs(brokerList.map(Integer.valueOf).toSeq.asJava) From b473a449317b4c0894c4032908db028b095893d0 Mon Sep 17 00:00:00 2001 From: wenbingshen Date: Sat, 13 Mar 2021 15:31:28 +0800 Subject: [PATCH 03/10] KAFKA-12454:Add ERROR logging on kafka-log-dirs when given brokerIds do not exist in current kafka cluster --- .../scala/kafka/admin/LogDirsCommand.scala | 25 +++++---- .../unit/kafka/admin/LogDirsCommandTest.scala | 52 +++++++++++++++++++ 2 files changed, 66 insertions(+), 11 deletions(-) create mode 100644 core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala diff --git a/core/src/main/scala/kafka/admin/LogDirsCommand.scala b/core/src/main/scala/kafka/admin/LogDirsCommand.scala index 46193ef3248f7..66e912186a9b4 100644 --- a/core/src/main/scala/kafka/admin/LogDirsCommand.scala +++ b/core/src/main/scala/kafka/admin/LogDirsCommand.scala @@ -41,23 +41,26 @@ object LogDirsCommand { val adminClient = createAdminClient(opts) val topicList = opts.options.valueOf(opts.topicListOpt).split(",").filter(!_.isEmpty) val clusterBrokers: Array[Int] = adminClient.describeCluster().nodes().get().asScala.map(_.id()).toArray + var nonExistBrokers: Array[Int] = Array.emptyIntArray val brokerList = Option(opts.options.valueOf(opts.brokerListOpt)) match { - case Some(brokerListStr) => brokerListStr.split(',').filter(!_.isEmpty).map(_.toInt) + case Some(brokerListStr) => { + val inputBrokers: Array[Int] = brokerListStr.split(',').filter(!_.isEmpty).map(_.toInt) + nonExistBrokers = inputBrokers.filterNot(brokerId => clusterBrokers.contains(brokerId)) + inputBrokers + } case None => clusterBrokers } - val nonExistBrokers: Array[Int] = brokerList.filterNot(brokerId => clusterBrokers.contains(brokerId)) if (!nonExistBrokers.isEmpty) { - System.err.println(s"The given node(s) does not exist from broker-list ${nonExistBrokers.mkString(",")}") - sys.exit(1) + out.println(s"ERROR: The given node(s) does not exist from broker-list ${nonExistBrokers.mkString(",")}") + } else { + out.println("Querying brokers for log directories information") + val describeLogDirsResult: DescribeLogDirsResult = adminClient.describeLogDirs(brokerList.map(Integer.valueOf).toSeq.asJava) + val logDirInfosByBroker = describeLogDirsResult.allDescriptions.get().asScala.map { case (k, v) => k -> v.asScala } + + out.println(s"Received log directory information from brokers ${brokerList.mkString(",")}") + out.println(formatAsJson(logDirInfosByBroker, topicList.toSet)) } - - out.println("Querying brokers for log directories information") - val describeLogDirsResult: DescribeLogDirsResult = adminClient.describeLogDirs(brokerList.map(Integer.valueOf).toSeq.asJava) - val logDirInfosByBroker = describeLogDirsResult.allDescriptions.get().asScala.map { case (k, v) => k -> v.asScala } - - out.println(s"Received log directory information from brokers ${brokerList.mkString(",")}") - out.println(formatAsJson(logDirInfosByBroker, topicList.toSet)) adminClient.close() } diff --git a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala new file mode 100644 index 0000000000000..aee1f4358d0c0 --- /dev/null +++ b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala @@ -0,0 +1,52 @@ +package unit.kafka.admin + +import java.io.{ByteArrayOutputStream, PrintStream} +import java.nio.charset.StandardCharsets + +import kafka.admin.LogDirsCommand +import kafka.integration.KafkaServerTestHarness +import kafka.server.KafkaConfig +import kafka.utils.TestUtils +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Test + +import scala.collection.Seq + +class LogDirsCommandTest extends KafkaServerTestHarness { + + def generateConfigs: Seq[KafkaConfig] = { + TestUtils.createBrokerConfigs(1, zkConnect) + .map(KafkaConfig.fromProps) + } + + @Test + def checkLogDirsCommandOutput(): Unit = { + val byteArrayOutputStream = new ByteArrayOutputStream + val printStream = new PrintStream(byteArrayOutputStream, false, StandardCharsets.UTF_8.name()) + //input exist brokerList + LogDirsCommand.describe(Array("--bootstrap-server", brokerList, "--broker-list", "0", "--describe"), printStream) + val existBrokersContent = new String(byteArrayOutputStream.toByteArray, StandardCharsets.UTF_8) + val existBrokersLineIter = existBrokersContent.split("\n").iterator + + assertTrue(existBrokersLineIter.hasNext) + assertTrue(existBrokersLineIter.next().contains(s"Querying brokers for log directories information")) + + //input nonExist brokerList + byteArrayOutputStream.reset() + LogDirsCommand.describe(Array("--bootstrap-server", brokerList, "--broker-list", "0,1,2", "--describe"), printStream) + val nonExistBrokersContent = new String(byteArrayOutputStream.toByteArray, StandardCharsets.UTF_8) + val nonExistBrokersLineIter = nonExistBrokersContent.split("\n").iterator + + assertTrue(nonExistBrokersLineIter.hasNext) + assertTrue(nonExistBrokersLineIter.next().contains(s"ERROR: The given node(s) does not exist from broker-list 1,2")) + + //use all brokerList for current cluster + byteArrayOutputStream.reset() + LogDirsCommand.describe(Array("--bootstrap-server", brokerList, "--describe"), printStream) + val allBrokersContent = new String(byteArrayOutputStream.toByteArray, StandardCharsets.UTF_8) + val allBrokersLineIter = allBrokersContent.split("\n").iterator + + assertTrue(allBrokersLineIter.hasNext) + assertTrue(allBrokersLineIter.next().contains(s"Querying brokers for log directories information")) + } +} From a4d6ae75b80e4a505ec8335cf2eaa87dc7db81ad Mon Sep 17 00:00:00 2001 From: wenbingshen Date: Mon, 15 Mar 2021 00:23:37 +0800 Subject: [PATCH 04/10] KAFKA-12454:Add ERROR logging on kafka-log-dirs when given brokerIds do not exist in current kafka cluster --- .../scala/kafka/admin/LogDirsCommand.scala | 40 ++++++++++--------- .../unit/kafka/admin/LogDirsCommandTest.scala | 2 +- 2 files changed, 22 insertions(+), 20 deletions(-) diff --git a/core/src/main/scala/kafka/admin/LogDirsCommand.scala b/core/src/main/scala/kafka/admin/LogDirsCommand.scala index 66e912186a9b4..b5b9679b0692d 100644 --- a/core/src/main/scala/kafka/admin/LogDirsCommand.scala +++ b/core/src/main/scala/kafka/admin/LogDirsCommand.scala @@ -39,29 +39,31 @@ object LogDirsCommand { def describe(args: Array[String], out: PrintStream): Unit = { val opts = new LogDirsCommandOptions(args) val adminClient = createAdminClient(opts) - val topicList = opts.options.valueOf(opts.topicListOpt).split(",").filter(!_.isEmpty) - val clusterBrokers: Array[Int] = adminClient.describeCluster().nodes().get().asScala.map(_.id()).toArray - var nonExistBrokers: Array[Int] = Array.emptyIntArray - val brokerList = Option(opts.options.valueOf(opts.brokerListOpt)) match { - case Some(brokerListStr) => { - val inputBrokers: Array[Int] = brokerListStr.split(',').filter(!_.isEmpty).map(_.toInt) - nonExistBrokers = inputBrokers.filterNot(brokerId => clusterBrokers.contains(brokerId)) - inputBrokers + val topicList = opts.options.valueOf(opts.topicListOpt).split(",").filter(_.nonEmpty) + var nonExistBrokers: Set[Int] = Set.empty + try { + val clusterBrokers: Set[Int] = adminClient.describeCluster().nodes().get().asScala.map(_.id()).toSet + val brokerList = Option(opts.options.valueOf(opts.brokerListOpt)) match { + case Some(brokerListStr) => + val inputBrokers: Set[Int] = brokerListStr.split(',').filter(_.nonEmpty).map(_.toInt).toSet + nonExistBrokers = inputBrokers.diff(clusterBrokers) + inputBrokers + case None => clusterBrokers } - case None => clusterBrokers - } - if (!nonExistBrokers.isEmpty) { - out.println(s"ERROR: The given node(s) does not exist from broker-list ${nonExistBrokers.mkString(",")}") - } else { - out.println("Querying brokers for log directories information") - val describeLogDirsResult: DescribeLogDirsResult = adminClient.describeLogDirs(brokerList.map(Integer.valueOf).toSeq.asJava) - val logDirInfosByBroker = describeLogDirsResult.allDescriptions.get().asScala.map { case (k, v) => k -> v.asScala } + if (nonExistBrokers.nonEmpty) { + out.println(s"ERROR: The given node(s) does not exist from broker-list: ${nonExistBrokers.mkString(",")}. Current cluster exist node(s): ${clusterBrokers.mkString(",")}") + } else { + out.println("Querying brokers for log directories information") + val describeLogDirsResult: DescribeLogDirsResult = adminClient.describeLogDirs(brokerList.map(Integer.valueOf).toSeq.asJava) + val logDirInfosByBroker = describeLogDirsResult.allDescriptions.get().asScala.map { case (k, v) => k -> v.asScala } - out.println(s"Received log directory information from brokers ${brokerList.mkString(",")}") - out.println(formatAsJson(logDirInfosByBroker, topicList.toSet)) + out.println(s"Received log directory information from brokers ${brokerList.mkString(",")}") + out.println(formatAsJson(logDirInfosByBroker, topicList.toSet)) + } + } finally { + adminClient.close() } - adminClient.close() } private def formatAsJson(logDirInfosByBroker: Map[Integer, Map[String, LogDirDescription]], topicSet: Set[String]): String = { diff --git a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala index aee1f4358d0c0..e7f70925291a8 100644 --- a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala +++ b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala @@ -38,7 +38,7 @@ class LogDirsCommandTest extends KafkaServerTestHarness { val nonExistBrokersLineIter = nonExistBrokersContent.split("\n").iterator assertTrue(nonExistBrokersLineIter.hasNext) - assertTrue(nonExistBrokersLineIter.next().contains(s"ERROR: The given node(s) does not exist from broker-list 1,2")) + assertTrue(nonExistBrokersLineIter.next().contains(s"ERROR: The given node(s) does not exist from broker-list: 1,2. Current cluster exist node(s): 0")) //use all brokerList for current cluster byteArrayOutputStream.reset() From fefdef90a38a80745bd1bf4d5beedf16f301c318 Mon Sep 17 00:00:00 2001 From: shenwenbing Date: Mon, 15 Mar 2021 22:57:02 +0800 Subject: [PATCH 05/10] KAFKA-12454:Add ERROR logging on kafka-log-dirs when given brokerIds do not exist in current kafka cluster --- .../scala/kafka/admin/LogDirsCommand.scala | 20 +++++++++---------- 1 file changed, 9 insertions(+), 11 deletions(-) diff --git a/core/src/main/scala/kafka/admin/LogDirsCommand.scala b/core/src/main/scala/kafka/admin/LogDirsCommand.scala index b5b9679b0692d..56d3990ec6222 100644 --- a/core/src/main/scala/kafka/admin/LogDirsCommand.scala +++ b/core/src/main/scala/kafka/admin/LogDirsCommand.scala @@ -40,25 +40,23 @@ object LogDirsCommand { val opts = new LogDirsCommandOptions(args) val adminClient = createAdminClient(opts) val topicList = opts.options.valueOf(opts.topicListOpt).split(",").filter(_.nonEmpty) - var nonExistBrokers: Set[Int] = Set.empty try { - val clusterBrokers: Set[Int] = adminClient.describeCluster().nodes().get().asScala.map(_.id()).toSet - val brokerList = Option(opts.options.valueOf(opts.brokerListOpt)) match { + val clusterBrokers = adminClient.describeCluster().nodes().get().asScala.map(_.id()).toSet + val (existingBrokers, nonExistingBrokers) = Option(opts.options.valueOf(opts.brokerListOpt)) match { case Some(brokerListStr) => - val inputBrokers: Set[Int] = brokerListStr.split(',').filter(_.nonEmpty).map(_.toInt).toSet - nonExistBrokers = inputBrokers.diff(clusterBrokers) - inputBrokers - case None => clusterBrokers + val inputBrokers = brokerListStr.split(',').filter(_.nonEmpty).map(_.toInt).toSet + (inputBrokers, inputBrokers.diff(clusterBrokers)) + case None => (clusterBrokers, Set.empty) } - if (nonExistBrokers.nonEmpty) { - out.println(s"ERROR: The given node(s) does not exist from broker-list: ${nonExistBrokers.mkString(",")}. Current cluster exist node(s): ${clusterBrokers.mkString(",")}") + if (nonExistingBrokers.nonEmpty) { + out.println(s"ERROR: The given node(s) does not exist from broker-list: ${nonExistingBrokers.mkString(",")}. Current cluster exist node(s): ${clusterBrokers.mkString(",")}") } else { out.println("Querying brokers for log directories information") - val describeLogDirsResult: DescribeLogDirsResult = adminClient.describeLogDirs(brokerList.map(Integer.valueOf).toSeq.asJava) + val describeLogDirsResult: DescribeLogDirsResult = adminClient.describeLogDirs(existingBrokers.map(Integer.valueOf).toSeq.asJava) val logDirInfosByBroker = describeLogDirsResult.allDescriptions.get().asScala.map { case (k, v) => k -> v.asScala } - out.println(s"Received log directory information from brokers ${brokerList.mkString(",")}") + out.println(s"Received log directory information from brokers ${existingBrokers.mkString(",")}") out.println(formatAsJson(logDirInfosByBroker, topicList.toSet)) } } finally { From c88b348516d299b0976c827cd0b86dbfc1a4cfe2 Mon Sep 17 00:00:00 2001 From: wenbingshen Date: Tue, 16 Mar 2021 00:15:08 +0800 Subject: [PATCH 06/10] KAFKA-12454:Add ERROR logging on kafka-log-dirs when given brokerIds do not exist in current kafka cluster --- .../unit/kafka/admin/LogDirsCommandTest.scala | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala index e7f70925291a8..56bf89151d26e 100644 --- a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala +++ b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala @@ -1,3 +1,19 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ package unit.kafka.admin import java.io.{ByteArrayOutputStream, PrintStream} From 848822ff5def3f27d0c265b8a8bd40fecc3e4dc4 Mon Sep 17 00:00:00 2001 From: wenbingshen Date: Tue, 16 Mar 2021 00:24:48 +0800 Subject: [PATCH 07/10] KAFKA-12454:Add ERROR logging on kafka-log-dirs when given brokerIds do not exist in current kafka cluster --- core/src/main/scala/kafka/admin/LogDirsCommand.scala | 6 +++--- .../test/scala/unit/kafka/admin/LogDirsCommandTest.scala | 2 +- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/core/src/main/scala/kafka/admin/LogDirsCommand.scala b/core/src/main/scala/kafka/admin/LogDirsCommand.scala index 56d3990ec6222..2bb3a97577213 100644 --- a/core/src/main/scala/kafka/admin/LogDirsCommand.scala +++ b/core/src/main/scala/kafka/admin/LogDirsCommand.scala @@ -21,7 +21,7 @@ import java.io.PrintStream import java.util.Properties import kafka.utils.{CommandDefaultOptions, CommandLineUtils, Json} -import org.apache.kafka.clients.admin.{Admin, AdminClientConfig, DescribeLogDirsResult, LogDirDescription} +import org.apache.kafka.clients.admin.{Admin, AdminClientConfig, LogDirDescription} import org.apache.kafka.common.utils.Utils import scala.jdk.CollectionConverters._ @@ -50,10 +50,10 @@ object LogDirsCommand { } if (nonExistingBrokers.nonEmpty) { - out.println(s"ERROR: The given node(s) does not exist from broker-list: ${nonExistingBrokers.mkString(",")}. Current cluster exist node(s): ${clusterBrokers.mkString(",")}") + out.println(s"ERROR: The given broker(s) does not exist from --broker-list: ${nonExistingBrokers.mkString(",")}. Current cluster exist broker(s): ${clusterBrokers.mkString(",")}") } else { out.println("Querying brokers for log directories information") - val describeLogDirsResult: DescribeLogDirsResult = adminClient.describeLogDirs(existingBrokers.map(Integer.valueOf).toSeq.asJava) + val describeLogDirsResult = adminClient.describeLogDirs(existingBrokers.map(Integer.valueOf).toSeq.asJava) val logDirInfosByBroker = describeLogDirsResult.allDescriptions.get().asScala.map { case (k, v) => k -> v.asScala } out.println(s"Received log directory information from brokers ${existingBrokers.mkString(",")}") diff --git a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala index 56bf89151d26e..437a7d67f8f07 100644 --- a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala +++ b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala @@ -54,7 +54,7 @@ class LogDirsCommandTest extends KafkaServerTestHarness { val nonExistBrokersLineIter = nonExistBrokersContent.split("\n").iterator assertTrue(nonExistBrokersLineIter.hasNext) - assertTrue(nonExistBrokersLineIter.next().contains(s"ERROR: The given node(s) does not exist from broker-list: 1,2. Current cluster exist node(s): 0")) + assertTrue(nonExistBrokersLineIter.next().contains(s"ERROR: The given broker(s) does not exist from --broker-list: 1,2. Current cluster exist broker(s): 0")) //use all brokerList for current cluster byteArrayOutputStream.reset() From 1e268327336deb2f5ef155c44c458b1b88a7c91b Mon Sep 17 00:00:00 2001 From: shenwenbing Date: Tue, 16 Mar 2021 10:04:56 +0800 Subject: [PATCH 08/10] KAFKA-12454:Add ERROR logging on kafka-log-dirs when given brokerIds do not exist in current kafka cluster --- core/src/main/scala/kafka/admin/LogDirsCommand.scala | 2 +- core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/admin/LogDirsCommand.scala b/core/src/main/scala/kafka/admin/LogDirsCommand.scala index 2bb3a97577213..1007a1bdcbd92 100644 --- a/core/src/main/scala/kafka/admin/LogDirsCommand.scala +++ b/core/src/main/scala/kafka/admin/LogDirsCommand.scala @@ -50,7 +50,7 @@ object LogDirsCommand { } if (nonExistingBrokers.nonEmpty) { - out.println(s"ERROR: The given broker(s) does not exist from --broker-list: ${nonExistingBrokers.mkString(",")}. Current cluster exist broker(s): ${clusterBrokers.mkString(",")}") + out.println(s"ERROR: The given brokers do not exist from --broker-list: ${nonExistingBrokers.mkString(",")}. Current cluster exist brokers: ${clusterBrokers.mkString(",")}") } else { out.println("Querying brokers for log directories information") val describeLogDirsResult = adminClient.describeLogDirs(existingBrokers.map(Integer.valueOf).toSeq.asJava) diff --git a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala index 437a7d67f8f07..697012d81357c 100644 --- a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala +++ b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala @@ -54,7 +54,7 @@ class LogDirsCommandTest extends KafkaServerTestHarness { val nonExistBrokersLineIter = nonExistBrokersContent.split("\n").iterator assertTrue(nonExistBrokersLineIter.hasNext) - assertTrue(nonExistBrokersLineIter.next().contains(s"ERROR: The given broker(s) does not exist from --broker-list: 1,2. Current cluster exist broker(s): 0")) + assertTrue(nonExistBrokersLineIter.next().contains(s"ERROR: The given brokers do not exist from --broker-list: 1,2. Current cluster exist brokers: 0")) //use all brokerList for current cluster byteArrayOutputStream.reset() From b245e5f64e03cb66963c44ac5ddb87cfc30e7825 Mon Sep 17 00:00:00 2001 From: shenwenbing Date: Tue, 16 Mar 2021 18:45:06 +0800 Subject: [PATCH 09/10] KAFKA-12454:Add ERROR logging on kafka-log-dirs when given brokerIds do not exist in current kafka cluster --- .../scala/kafka/admin/LogDirsCommand.scala | 5 ++-- .../unit/kafka/admin/LogDirsCommandTest.scala | 27 ++++++++++++------- 2 files changed, 21 insertions(+), 11 deletions(-) diff --git a/core/src/main/scala/kafka/admin/LogDirsCommand.scala b/core/src/main/scala/kafka/admin/LogDirsCommand.scala index 1007a1bdcbd92..a15e4eb5225bc 100644 --- a/core/src/main/scala/kafka/admin/LogDirsCommand.scala +++ b/core/src/main/scala/kafka/admin/LogDirsCommand.scala @@ -39,8 +39,8 @@ object LogDirsCommand { def describe(args: Array[String], out: PrintStream): Unit = { val opts = new LogDirsCommandOptions(args) val adminClient = createAdminClient(opts) - val topicList = opts.options.valueOf(opts.topicListOpt).split(",").filter(_.nonEmpty) try { + val topicList = opts.options.valueOf(opts.topicListOpt).split(",").filter(_.nonEmpty) val clusterBrokers = adminClient.describeCluster().nodes().get().asScala.map(_.id()).toSet val (existingBrokers, nonExistingBrokers) = Option(opts.options.valueOf(opts.brokerListOpt)) match { case Some(brokerListStr) => @@ -50,7 +50,8 @@ object LogDirsCommand { } if (nonExistingBrokers.nonEmpty) { - out.println(s"ERROR: The given brokers do not exist from --broker-list: ${nonExistingBrokers.mkString(",")}. Current cluster exist brokers: ${clusterBrokers.mkString(",")}") + out.println(s"ERROR: The given brokers do not exist from --broker-list: ${nonExistingBrokers.mkString(",")}." + + s" Current existent brokers: ${clusterBrokers.mkString(",")}") } else { out.println("Querying brokers for log directories information") val describeLogDirsResult = adminClient.describeLogDirs(existingBrokers.map(Integer.valueOf).toSeq.asJava) diff --git a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala index 697012d81357c..43889902980dc 100644 --- a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala +++ b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala @@ -41,20 +41,29 @@ class LogDirsCommandTest extends KafkaServerTestHarness { val printStream = new PrintStream(byteArrayOutputStream, false, StandardCharsets.UTF_8.name()) //input exist brokerList LogDirsCommand.describe(Array("--bootstrap-server", brokerList, "--broker-list", "0", "--describe"), printStream) - val existBrokersContent = new String(byteArrayOutputStream.toByteArray, StandardCharsets.UTF_8) - val existBrokersLineIter = existBrokersContent.split("\n").iterator + val existingBrokersContent = new String(byteArrayOutputStream.toByteArray, StandardCharsets.UTF_8) + val existingBrokersLineIter = existingBrokersContent.split("\n").iterator - assertTrue(existBrokersLineIter.hasNext) - assertTrue(existBrokersLineIter.next().contains(s"Querying brokers for log directories information")) + assertTrue(existingBrokersLineIter.hasNext) + assertTrue(existingBrokersLineIter.next().contains(s"Querying brokers for log directories information")) - //input nonExist brokerList + //input nonexistent brokerList byteArrayOutputStream.reset() LogDirsCommand.describe(Array("--bootstrap-server", brokerList, "--broker-list", "0,1,2", "--describe"), printStream) - val nonExistBrokersContent = new String(byteArrayOutputStream.toByteArray, StandardCharsets.UTF_8) - val nonExistBrokersLineIter = nonExistBrokersContent.split("\n").iterator + val nonExistingBrokersContent = new String(byteArrayOutputStream.toByteArray, StandardCharsets.UTF_8) + val nonExistingBrokersLineIter = nonExistingBrokersContent.split("\n").iterator - assertTrue(nonExistBrokersLineIter.hasNext) - assertTrue(nonExistBrokersLineIter.next().contains(s"ERROR: The given brokers do not exist from --broker-list: 1,2. Current cluster exist brokers: 0")) + assertTrue(nonExistingBrokersLineIter.hasNext) + assertTrue(nonExistingBrokersLineIter.next().contains(s"ERROR: The given brokers do not exist from --broker-list: 1,2. Current existent brokers: 0")) + + //input duplicate ids + byteArrayOutputStream.reset() + LogDirsCommand.describe(Array("--bootstrap-server", brokerList, "--broker-list", "0,0,1,2,2", "--describe"), printStream) + val duplicateBrokersContent = new String(byteArrayOutputStream.toByteArray, StandardCharsets.UTF_8) + val duplicateBrokersLineIter = duplicateBrokersContent.split("\n").iterator + + assertTrue(duplicateBrokersLineIter.hasNext) + assertTrue(duplicateBrokersLineIter.next().contains(s"ERROR: The given brokers do not exist from --broker-list: 1,2. Current existent brokers: 0")) //use all brokerList for current cluster byteArrayOutputStream.reset() From 24f520febb060c90918b72542a91574b900367ea Mon Sep 17 00:00:00 2001 From: shenwenbing Date: Wed, 17 Mar 2021 14:29:54 +0800 Subject: [PATCH 10/10] KAFKA-12454:Add ERROR logging on kafka-log-dirs when given brokerIds do not exist in current kafka cluster --- core/src/main/scala/kafka/admin/LogDirsCommand.scala | 2 +- core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala | 3 +-- 2 files changed, 2 insertions(+), 3 deletions(-) diff --git a/core/src/main/scala/kafka/admin/LogDirsCommand.scala b/core/src/main/scala/kafka/admin/LogDirsCommand.scala index a15e4eb5225bc..d8c802e7d0e23 100644 --- a/core/src/main/scala/kafka/admin/LogDirsCommand.scala +++ b/core/src/main/scala/kafka/admin/LogDirsCommand.scala @@ -45,7 +45,7 @@ object LogDirsCommand { val (existingBrokers, nonExistingBrokers) = Option(opts.options.valueOf(opts.brokerListOpt)) match { case Some(brokerListStr) => val inputBrokers = brokerListStr.split(',').filter(_.nonEmpty).map(_.toInt).toSet - (inputBrokers, inputBrokers.diff(clusterBrokers)) + (inputBrokers.intersect(clusterBrokers), inputBrokers.diff(clusterBrokers)) case None => (clusterBrokers, Set.empty) } diff --git a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala index 43889902980dc..397d6d580c798 100644 --- a/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala +++ b/core/src/test/scala/unit/kafka/admin/LogDirsCommandTest.scala @@ -14,12 +14,11 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package unit.kafka.admin +package kafka.admin import java.io.{ByteArrayOutputStream, PrintStream} import java.nio.charset.StandardCharsets -import kafka.admin.LogDirsCommand import kafka.integration.KafkaServerTestHarness import kafka.server.KafkaConfig import kafka.utils.TestUtils