From be960b81ef7e8ba195ffaf665eac23186019b5c8 Mon Sep 17 00:00:00 2001 From: Chia-Ping Tsai Date: Sat, 13 Jun 2020 13:37:21 +0800 Subject: [PATCH] HOTFIX: fix compile error in o.a.k.c.u.TopicAdminTest --- .../kafka/connect/util/TopicAdminTest.java | 33 ++++++++++--------- 1 file changed, 18 insertions(+), 15 deletions(-) diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/util/TopicAdminTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/util/TopicAdminTest.java index b655664c4b181..dd3062e4159bc 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/util/TopicAdminTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/util/TopicAdminTest.java @@ -36,6 +36,7 @@ import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.message.CreateTopicsResponseData; import org.apache.kafka.common.message.CreateTopicsResponseData.CreatableTopicResult; +import org.apache.kafka.common.message.DescribeConfigsResponseData; import org.apache.kafka.common.message.MetadataResponseData; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.requests.ApiError; @@ -46,13 +47,14 @@ import org.apache.kafka.connect.errors.ConnectException; import org.junit.Test; -import java.util.ArrayList; -import java.util.Collection; import java.util.Collections; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.ExecutionException; +import java.util.stream.Collectors; +import java.util.stream.Stream; import static org.apache.kafka.common.message.MetadataResponseData.MetadataResponseTopic; import static org.junit.Assert.assertEquals; @@ -567,19 +569,20 @@ private DescribeConfigsResponse describeConfigsResponseWithTopicAuthorizationExc } private DescribeConfigsResponse describeConfigsResponse(ApiError error, NewTopic... topics) { - if (error == null) error = new ApiError(Errors.NONE, ""); - Map configs = new HashMap<>(); - for (NewTopic topic : topics) { - ConfigResource resource = new ConfigResource(ConfigResource.Type.TOPIC, topic.name()); - DescribeConfigsResponse.ConfigSource source = DescribeConfigsResponse.ConfigSource.TOPIC_CONFIG; - Collection entries = new ArrayList<>(); - topic.configs().forEach((k, v) -> new DescribeConfigsResponse.ConfigEntry( - k, v, source, false, false, Collections.emptySet() - )); - DescribeConfigsResponse.Config config = new DescribeConfigsResponse.Config(error, entries); - configs.put(resource, config); - } - return new DescribeConfigsResponse(1000, configs); + List results = Stream.of(topics) + .map(topic -> new DescribeConfigsResponseData.DescribeConfigsResult() + .setErrorCode(error.error().code()) + .setErrorMessage(error.message()) + .setResourceType(ConfigResource.Type.TOPIC.id()) + .setResourceName(topic.name()) + .setConfigs(topic.configs().entrySet() + .stream() + .map(e -> new DescribeConfigsResponseData.DescribeConfigsResourceResult() + .setName(e.getKey()) + .setValue(e.getValue())) + .collect(Collectors.toList()))) + .collect(Collectors.toList()); + return new DescribeConfigsResponse(new DescribeConfigsResponseData().setThrottleTimeMs(1000).setResults(results)); } }