From 90308c884b7e88a876ddcd5013c9a2234ded48f9 Mon Sep 17 00:00:00 2001 From: Manikumar Reddy Date: Tue, 21 Jul 2020 01:07:48 +0530 Subject: [PATCH] KAFKA-9432:(follow-up) Set configKeys to null in describeConfigs() to make it backward compatible with older Kafka versions. --- .../kafka/clients/admin/KafkaAdminClient.java | 3 +- .../unit/kafka/server/AdminManagerTest.scala | 20 +++++++++- .../client_compatibility_features_test.py | 2 + .../kafka/tools/ClientCompatibilityTest.java | 37 +++++++++++++++++++ 4 files changed, 60 insertions(+), 2 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 8bd8142aff823..99058abe3388d 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -1949,7 +1949,8 @@ DescribeConfigsRequest.Builder createRequest(int timeoutMs) { .map(config -> new DescribeConfigsRequestData.DescribeConfigsResource() .setResourceName(config.name()) - .setResourceType(config.type().id())) + .setResourceType(config.type().id()) + .setConfigurationKeys(null)) .collect(Collectors.toList())) .setIncludeSynonyms(options.includeSynonyms()) .setIncludeDocumentation(options.includeDocumentation())); diff --git a/core/src/test/scala/unit/kafka/server/AdminManagerTest.scala b/core/src/test/scala/unit/kafka/server/AdminManagerTest.scala index b472bcfb8d09b..ad0efd6aafaa4 100644 --- a/core/src/test/scala/unit/kafka/server/AdminManagerTest.scala +++ b/core/src/test/scala/unit/kafka/server/AdminManagerTest.scala @@ -28,6 +28,7 @@ import org.apache.kafka.common.protocol.Errors import org.junit.{After, Test} import org.junit.Assert.assertEquals +import org.junit.Assert.assertFalse class AdminManagerTest { @@ -60,6 +61,23 @@ class AdminManagerTest { .setConfigurationKeys(null)) val adminManager = createAdminManager() val results: List[DescribeConfigsResponseData.DescribeConfigsResult] = adminManager.describeConfigs(resources, true, true) - assertEquals(results.head.errorCode(), Errors.NONE.code) + assertEquals(Errors.NONE.code, results.head.errorCode()) + assertFalse("Should return configs", results.head.configs().isEmpty) + } + + @Test + def testDescribeConfigsWithEmptyConfigurationKeys(): Unit = { + EasyMock.expect(zkClient.getEntityConfigs(ConfigType.Topic, topic)).andReturn(TestUtils.createBrokerConfig(brokerId, "zk")) + EasyMock.expect(metadataCache.contains(topic)).andReturn(true) + + EasyMock.replay(zkClient, metadataCache) + + val resources = List(new DescribeConfigsRequestData.DescribeConfigsResource() + .setResourceName(topic) + .setResourceType(ConfigResource.Type.TOPIC.id)) + val adminManager = createAdminManager() + val results: List[DescribeConfigsResponseData.DescribeConfigsResult] = adminManager.describeConfigs(resources, true, true) + assertEquals(Errors.NONE.code, results.head.errorCode()) + assertFalse("Should return configs", results.head.configs().isEmpty) } } diff --git a/tests/kafkatest/tests/client/client_compatibility_features_test.py b/tests/kafkatest/tests/client/client_compatibility_features_test.py index 7a0d471398b0f..50f4f7f8af69b 100644 --- a/tests/kafkatest/tests/client/client_compatibility_features_test.py +++ b/tests/kafkatest/tests/client/client_compatibility_features_test.py @@ -39,8 +39,10 @@ def get_broker_features(broker_version): features["expect-record-too-large-exception"] = False if broker_version < V_0_11_0_0: features["describe-acls-supported"] = False + features["describe-configs-supported"] = False else: features["describe-acls-supported"] = True + features["describe-configs-supported"] = True return features def run_command(node, cmd, ssh_log_file): diff --git a/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java b/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java index 89e6b44c8d2be..a5d6c7a835dbb 100644 --- a/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java +++ b/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java @@ -22,6 +22,7 @@ import net.sourceforge.argparse4j.inf.Namespace; import org.apache.kafka.clients.admin.Admin; import org.apache.kafka.clients.admin.AdminClientConfig; +import org.apache.kafka.clients.admin.Config; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.admin.TopicListing; import org.apache.kafka.clients.consumer.ConsumerConfig; @@ -40,6 +41,7 @@ import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.acl.AclBindingFilter; +import org.apache.kafka.common.config.ConfigResource; import org.apache.kafka.common.errors.RecordTooLargeException; import org.apache.kafka.common.errors.SecurityDisabledException; import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; @@ -85,6 +87,7 @@ static class TestConfig { final int numClusterNodes; final boolean createTopicsSupported; final boolean describeAclsSupported; + final boolean describeConfigsSupported; TestConfig(Namespace res) { this.bootstrapServer = res.getString("bootstrapServer"); @@ -95,6 +98,7 @@ static class TestConfig { this.numClusterNodes = res.getInt("numClusterNodes"); this.createTopicsSupported = res.getBoolean("createTopicsSupported"); this.describeAclsSupported = res.getBoolean("describeAclsSupported"); + this.describeConfigsSupported = res.getBoolean("describeConfigsSupported"); } } @@ -161,6 +165,13 @@ public static void main(String[] args) throws Exception { .dest("describeAclsSupported") .metavar("DESCRIBE_ACLS_SUPPORTED") .help("Whether describeAcls is supported in the AdminClient."); + parser.addArgument("--describe-configs-supported") + .action(store()) + .required(true) + .type(Boolean.class) + .dest("describeConfigsSupported") + .metavar("DESCRIBE_CONFIGS_SUPPORTED") + .help("Whether describeConfigs is supported in the AdminClient."); Namespace res = null; try { @@ -260,6 +271,9 @@ void testAdminClient() throws Throwable { log.info("Saw only {} cluster nodes. Waiting to see {}.", nodes.size(), testConfig.numClusterNodes); } + + testDescribeConfigsMethod(client); + tryFeature("createTopics", testConfig.createTopicsSupported, () -> { try { @@ -297,6 +311,29 @@ void testAdminClient() throws Throwable { } } + private void testDescribeConfigsMethod(final Admin client) throws Throwable { + tryFeature("describeConfigsSupported", testConfig.describeConfigsSupported, + () -> { + try { + Collection nodes = client.describeCluster().nodes().get(); + + final ConfigResource configResource = new ConfigResource( + ConfigResource.Type.BROKER, + nodes.iterator().next().idString() + ); + + Map brokerConfig = + client.describeConfigs(Collections.singleton(configResource)).all().get(); + + if (brokerConfig.get(configResource).entries().isEmpty()) { + throw new KafkaException("Expected to see config entries, but got zero entries"); + } + } catch (ExecutionException e) { + throw e.getCause(); + } + }); + } + private void createTopicsResultTest(Admin client, Collection topics) throws InterruptedException, ExecutionException { while (true) {