From 8b2520da1070e0067e43881ea9075aa85baff793 Mon Sep 17 00:00:00 2001 From: "Joseph (Ting-Chou) Lin" Date: Wed, 19 Jan 2022 02:01:21 +0800 Subject: [PATCH 1/3] [LI-HOTFIX] Add test case for describe maintenance broker config TICKET = KAFKA-8527 LI_DESCRIPTION = In the commit 6b93c274b9 `[LI-HOTFIX] Add dynamic maintenance broker config`, the test cases didn't cover the describe config API. Adding this to protect the code path and facilitate the future rebase process. EXIT_CRITERIA = TICKET [KAFKA-8527] or the patch is dropped --- .../api/PlaintextAdminIntegrationTest.scala | 20 +++++++++++++++- .../tools/MaintenanceBrokerTestUtils.scala | 24 +++++++++++++++++++ .../kafka/server/MaintenanceBrokerTest.scala | 12 ++-------- 3 files changed, 45 insertions(+), 11 deletions(-) create mode 100644 core/src/test/scala/integration/kafka/tools/MaintenanceBrokerTestUtils.scala diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index b96a1024a6c73..3902d63d4de38 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -26,9 +26,10 @@ import java.util.{Collections, Optional, Properties} import java.{time, util} import integration.kafka.api.BaseAdminIntegrationTest +import integration.kafka.tools.MaintenanceBrokerTestUtils import kafka.log.LogConfig import kafka.security.auth.Group -import kafka.server.{Defaults, KafkaConfig, KafkaServer} +import kafka.server.{Defaults, DynamicConfig, KafkaConfig, KafkaServer} import kafka.utils.TestUtils._ import kafka.utils.{Log4jController, TestUtils} import kafka.zk.KafkaZkClient @@ -891,6 +892,23 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertTrue(intercept[ExecutionException](describeResult2.values.get(invalidTopic).get).getCause.isInstanceOf[InvalidTopicException]) } + @Test + def testDescribeConfigsForMaintenanceBroker(): Unit = { + client = AdminClient.create(createConfig) + val maintenanceBrokerIds = Seq(0, 2) + MaintenanceBrokerTestUtils.setMaintenanceBrokers(adminZkClient, zkClient, servers, maintenanceBrokerIds) + servers.map(_.config.brokerId).foreach { brokerId => + val describedMaintenanceBrokerConfigStr = + client.describeConfigs(Collections.singleton(new ConfigResource(ConfigResource.Type.BROKER, brokerId.toString))) + .all().get(5, TimeUnit.SECONDS) + .values().asScala.head.get(DynamicConfig.Broker.MaintenanceBrokerListProp) + .value() + assertNotNull(describedMaintenanceBrokerConfigStr) + val describedMaintenanceBrokerIds = describedMaintenanceBrokerConfigStr.split(",").map(_.toInt).toSet + assertEquals(describedMaintenanceBrokerIds, maintenanceBrokerIds.toSet) + } + } + private def subscribeAndWaitForAssignment(topic: String, consumer: KafkaConsumer[Array[Byte], Array[Byte]]): Unit = { consumer.subscribe(Collections.singletonList(topic)) TestUtils.pollUntilTrue(consumer, () => !consumer.assignment.isEmpty, "Expected non-empty assignment") diff --git a/core/src/test/scala/integration/kafka/tools/MaintenanceBrokerTestUtils.scala b/core/src/test/scala/integration/kafka/tools/MaintenanceBrokerTestUtils.scala new file mode 100644 index 0000000000000..87ae13f5cc10e --- /dev/null +++ b/core/src/test/scala/integration/kafka/tools/MaintenanceBrokerTestUtils.scala @@ -0,0 +1,24 @@ +package integration.kafka.tools + +import kafka.server.{DynamicConfig, KafkaServer} +import kafka.utils.CoreUtils.propsWith +import kafka.utils.TestUtils +import kafka.zk.{AdminZkClient, KafkaZkClient} + +object MaintenanceBrokerTestUtils { + + def setMaintenanceBrokers(adminZkClient: AdminZkClient, + zkClient: KafkaZkClient, + brokers: Seq[KafkaServer], + maintenanceBrokerIds: Seq[Int]): Unit = { + val propstring = maintenanceBrokerIds.mkString(",") + adminZkClient.changeBrokerConfig(None, + propsWith((DynamicConfig.Broker.MaintenanceBrokerListProp, propstring))) + + val controllerId = TestUtils.waitUntilControllerElected(zkClient) + + TestUtils.waitUntilTrue(() => brokers(controllerId).config.getMaintenanceBrokerList == maintenanceBrokerIds, + s"wait until broker $propstring is masked as maintenance broker not taking new partitions", 5000) + } + +} diff --git a/core/src/test/scala/unit/kafka/server/MaintenanceBrokerTest.scala b/core/src/test/scala/unit/kafka/server/MaintenanceBrokerTest.scala index c97740136623b..0f7586d98cc55 100644 --- a/core/src/test/scala/unit/kafka/server/MaintenanceBrokerTest.scala +++ b/core/src/test/scala/unit/kafka/server/MaintenanceBrokerTest.scala @@ -19,6 +19,7 @@ package kafka.server import java.util.{Optional, Properties} +import integration.kafka.tools.MaintenanceBrokerTestUtils import kafka.server.KafkaConfig.fromProps import kafka.utils.CoreUtils._ import kafka.utils.TestUtils @@ -201,15 +202,6 @@ class MaintenanceBrokerTest extends ZooKeeperTestHarness { } def setMaintenanceBrokers(brokerIds: Seq[Int]): Unit = { - var propstring = brokerIds.mkString(",") - adminZkClient.changeBrokerConfig(None, - propsWith((DynamicConfig.Broker.MaintenanceBrokerListProp, propstring))) - - val controllerId = TestUtils.waitUntilControllerElected(zkClient) - - TestUtils.waitUntilTrue(() => brokers(controllerId).config.getMaintenanceBrokerList == brokerIds, - s"wait until broker $propstring is masked as maintenance broker not taking new partitions", 5000) - + MaintenanceBrokerTestUtils.setMaintenanceBrokers(adminZkClient, zkClient, brokers, brokerIds) } - } From 9f0a9a460725dbcc0cd33df2a35694d854769eb2 Mon Sep 17 00:00:00 2001 From: "Joseph (Ting-Chou) Lin" Date: Wed, 19 Jan 2022 02:46:09 +0800 Subject: [PATCH 2/3] Add license --- .../tools/MaintenanceBrokerTestUtils.scala | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/core/src/test/scala/integration/kafka/tools/MaintenanceBrokerTestUtils.scala b/core/src/test/scala/integration/kafka/tools/MaintenanceBrokerTestUtils.scala index 87ae13f5cc10e..82edae2829139 100644 --- a/core/src/test/scala/integration/kafka/tools/MaintenanceBrokerTestUtils.scala +++ b/core/src/test/scala/integration/kafka/tools/MaintenanceBrokerTestUtils.scala @@ -1,3 +1,20 @@ +/** + * 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 integration.kafka.tools import kafka.server.{DynamicConfig, KafkaServer} From a77ed625b0ce96754c0e8cce57d61df59f3953a9 Mon Sep 17 00:00:00 2001 From: "Joseph (Ting-Chou) Lin" Date: Wed, 19 Jan 2022 02:47:19 +0800 Subject: [PATCH 3/3] Fix some code style nits --- .../scala/unit/kafka/server/MaintenanceBrokerTest.scala | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/core/src/test/scala/unit/kafka/server/MaintenanceBrokerTest.scala b/core/src/test/scala/unit/kafka/server/MaintenanceBrokerTest.scala index 0f7586d98cc55..73d4df8d9a2f4 100644 --- a/core/src/test/scala/unit/kafka/server/MaintenanceBrokerTest.scala +++ b/core/src/test/scala/unit/kafka/server/MaintenanceBrokerTest.scala @@ -17,22 +17,20 @@ package kafka.server -import java.util.{Optional, Properties} +import java.util.Properties import integration.kafka.tools.MaintenanceBrokerTestUtils import kafka.server.KafkaConfig.fromProps -import kafka.utils.CoreUtils._ import kafka.utils.TestUtils import kafka.utils.TestUtils._ import kafka.zk.ZooKeeperTestHarness import org.apache.kafka.clients.admin._ import org.apache.kafka.common.network.ListenerName import org.apache.kafka.common.security.auth.SecurityProtocol - -import scala.collection.JavaConverters._ import org.junit.Assert._ import org.junit.{After, Test} +import scala.collection.JavaConverters._ import scala.collection.Map /** @@ -204,4 +202,5 @@ class MaintenanceBrokerTest extends ZooKeeperTestHarness { def setMaintenanceBrokers(brokerIds: Seq[Int]): Unit = { MaintenanceBrokerTestUtils.setMaintenanceBrokers(adminZkClient, zkClient, brokers, brokerIds) } + }