From 56fe71bb497cf41675d717583062818a391b3e34 Mon Sep 17 00:00:00 2001 From: Chia-Ping Tsai Date: Fri, 19 Feb 2021 02:52:50 +0800 Subject: [PATCH] KAFKA-12339 Starting new connector cluster with new internal topics encounters UnknownTopicOrPartitionException --- .../internals/MetadataOperationContext.java | 1 + .../clients/admin/KafkaAdminClientTest.java | 54 ++++++++++++++++--- 2 files changed, 47 insertions(+), 8 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/internals/MetadataOperationContext.java b/clients/src/main/java/org/apache/kafka/clients/admin/internals/MetadataOperationContext.java index c05e5cfac0f3c..e7f2c07d9de83 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/internals/MetadataOperationContext.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/internals/MetadataOperationContext.java @@ -82,6 +82,7 @@ public Collection topics() { public static void handleMetadataErrors(MetadataResponse response) { for (TopicMetadata tm : response.topicMetadata()) { + if (shouldRefreshMetadata(tm.error())) throw tm.error().exception(); for (PartitionMetadata pm : tm.partitionMetadata()) { if (shouldRefreshMetadata(pm.error)) { throw pm.error.exception(); diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index ec05d2c002fd1..df51949e1eea0 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -458,12 +458,16 @@ private static FindCoordinatorResponse prepareFindCoordinatorResponse(Errors err } private static MetadataResponse prepareMetadataResponse(Cluster cluster, Errors error) { + return prepareMetadataResponse(cluster, error, error); + } + + private static MetadataResponse prepareMetadataResponse(Cluster cluster, Errors topicError, Errors partitionError) { List metadata = new ArrayList<>(); for (String topic : cluster.topics()) { List pms = new ArrayList<>(); for (PartitionInfo pInfo : cluster.availablePartitionsForTopic(topic)) { MetadataResponsePartition pm = new MetadataResponsePartition() - .setErrorCode(error.code()) + .setErrorCode(partitionError.code()) .setPartitionIndex(pInfo.partition()) .setLeaderId(pInfo.leader().id()) .setLeaderEpoch(234) @@ -473,19 +477,19 @@ private static MetadataResponse prepareMetadataResponse(Cluster cluster, Errors pms.add(pm); } MetadataResponseTopic tm = new MetadataResponseTopic() - .setErrorCode(error.code()) + .setErrorCode(topicError.code()) .setName(topic) .setIsInternal(false) .setPartitions(pms); metadata.add(tm); } return MetadataResponse.prepareResponse(true, - 0, - cluster.nodes(), - cluster.clusterResource().clusterId(), - cluster.controller().id(), - metadata, - MetadataResponse.AUTHORIZED_OPERATIONS_OMITTED); + 0, + cluster.nodes(), + cluster.clusterResource().clusterId(), + cluster.controller().id(), + metadata, + MetadataResponse.AUTHORIZED_OPERATIONS_OMITTED); } private static DescribeGroupsResponseData prepareDescribeGroupsResponseData(String groupId, @@ -4058,6 +4062,40 @@ public void testListOffsets() throws Exception { } } + @Test + public void testListOffsetsRetriableErrorOnMetadata() throws Exception { + Node node = new Node(0, "localhost", 8120); + List nodes = Collections.singletonList(node); + final Cluster cluster = new Cluster( + "mockClusterId", + nodes, + Collections.singleton(new PartitionInfo("foo", 0, node, new Node[]{node}, new Node[]{node})), + Collections.emptySet(), + Collections.emptySet(), + node); + final TopicPartition tp0 = new TopicPartition("foo", 0); + + try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(cluster)) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + env.kafkaClient().prepareResponse(prepareMetadataResponse(cluster, Errors.UNKNOWN_TOPIC_OR_PARTITION, Errors.NONE)); + // metadata refresh because of UNKNOWN_TOPIC_OR_PARTITION + env.kafkaClient().prepareResponse(prepareMetadataResponse(cluster, Errors.NONE)); + // listoffsets response from broker 0 + ListOffsetsResponseData responseData = new ListOffsetsResponseData() + .setThrottleTimeMs(0) + .setTopics(Collections.singletonList(ListOffsetsResponse.singletonListOffsetsTopicResponse(tp0, Errors.NONE, -1L, 123L, 321))); + env.kafkaClient().prepareResponseFrom(new ListOffsetsResponse(responseData), node); + + ListOffsetsResult result = env.adminClient().listOffsets(Collections.singletonMap(tp0, OffsetSpec.latest())); + + Map offsets = result.all().get(3, TimeUnit.SECONDS); + assertEquals(1, offsets.size()); + assertEquals(123L, offsets.get(tp0).offset()); + assertEquals(321, offsets.get(tp0).leaderEpoch().get().intValue()); + assertEquals(-1L, offsets.get(tp0).timestamp()); + } + } + @Test public void testListOffsetsRetriableErrors() throws Exception {