CalvinConfluent commented on code in PR #15470:
URL: https://github.com/apache/kafka/pull/15470#discussion_r1515288246
##########
clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java:
##########
@@ -2184,9 +2177,122 @@ void handleFailure(Throwable throwable) {
completeAllExceptionally(topicFutures.values(), throwable);
}
};
+ return call;
+ }
+
+ @SuppressWarnings("MethodLength")
+ private Map<String, KafkaFuture<TopicDescription>>
handleDescribeTopicsByNamesWithDescribeTopicPartitionsApi(final
Collection<String> topicNames, DescribeTopicsOptions options) {
+ final Map<String, KafkaFutureImpl<TopicDescription>> topicFutures =
new HashMap<>(topicNames.size());
+ final ArrayList<String> topicNamesList = new ArrayList<>();
+ for (String topicName : topicNames) {
+ if (topicNameIsUnrepresentable(topicName)) {
+ KafkaFutureImpl<TopicDescription> future = new
KafkaFutureImpl<>();
+ future.completeExceptionally(new InvalidTopicException("The
given topic name '" +
+ topicName + "' cannot be represented in a request."));
+ topicFutures.put(topicName, future);
+ } else if (!topicFutures.containsKey(topicName)) {
+ topicFutures.put(topicName, new KafkaFutureImpl<>());
+ topicNamesList.add(topicName);
+ }
+ }
+ final long now = time.milliseconds();
+ Call call = new Call("describeTopicPartitions", calcDeadlineMs(now,
options.timeoutMs()),
+ new LeastLoadedNodeProvider()) {
+ Map<String, TopicRequest> pendingTopics =
+ topicNamesList.stream().map(topicName -> new
TopicRequest().setName(topicName))
+ .collect(Collectors.toMap(topicRequest ->
topicRequest.name(), topicRequest -> topicRequest, (t1, t2) -> t1,
TreeMap::new));
+
+ DescribeTopicPartitionsRequestData.Cursor requestCursor = null;
+ TopicDescription partiallyFinishedTopicDescription = null;
+
+ @Override
+ DescribeTopicPartitionsRequest.Builder createRequest(int
timeoutMs) {
+ DescribeTopicPartitionsRequestData request = new
DescribeTopicPartitionsRequestData()
+ .setTopics(new ArrayList<>(pendingTopics.values()))
+
.setResponsePartitionLimit(options.partitionSizeLimitPerResponse());
+ request.setCursor(requestCursor);
+ return new DescribeTopicPartitionsRequest.Builder(request);
+ }
+
+ @Override
+ void handleResponse(AbstractResponse abstractResponse) {
+ DescribeTopicPartitionsResponse response =
(DescribeTopicPartitionsResponse) abstractResponse;
+ DescribeTopicPartitionsResponseData.Cursor responseCursor =
response.data().nextCursor();
+
+ for (DescribeTopicPartitionsResponseTopic topic :
response.data().topics()) {
+ String topicName = topic.name();
+ Errors error = Errors.forCode(topic.errorCode());
+
+ KafkaFutureImpl<TopicDescription> future =
topicFutures.get(topicName);
+ if (error != Errors.NONE) {
+ future.completeExceptionally(error.exception());
+ topicFutures.remove(topicName);
Review Comment:
Yes, we don't need to remove it.
Also added a failure handling test.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]