From 909da08a6d8f408dcdcb64ea4313786eba2be650 Mon Sep 17 00:00:00 2001 From: Andrew Schofield Date: Tue, 5 Mar 2024 19:42:09 +0000 Subject: [PATCH 1/2] KAFKA-16319: Divide DeleteTopics requests by leader node --- .../admin/internals/DeleteRecordsHandler.java | 6 +- .../internals/DeleteRecordsHandlerTest.java | 67 +++++++++++++++++-- 2 files changed, 66 insertions(+), 7 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/internals/DeleteRecordsHandler.java b/clients/src/main/java/org/apache/kafka/clients/admin/internals/DeleteRecordsHandler.java index 2daad226034d2..9f40d19f00b48 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/internals/DeleteRecordsHandler.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/internals/DeleteRecordsHandler.java @@ -79,15 +79,15 @@ public static SimpleAdminApiFuture newFuture( @Override public DeleteRecordsRequest.Builder buildBatchedRequest(int brokerId, Set keys) { Map deletionsForTopic = new HashMap<>(); - for (Map.Entry entry: recordsToDelete.entrySet()) { - TopicPartition topicPartition = entry.getKey(); + for (TopicPartition topicPartition : keys) { + RecordsToDelete toDelete = recordsToDelete.get(topicPartition); DeleteRecordsRequestData.DeleteRecordsTopic deleteRecords = deletionsForTopic.computeIfAbsent( topicPartition.topic(), key -> new DeleteRecordsRequestData.DeleteRecordsTopic().setName(topicPartition.topic()) ); deleteRecords.partitions().add(new DeleteRecordsRequestData.DeleteRecordsPartition() .setPartitionIndex(topicPartition.partition()) - .setOffset(entry.getValue().beforeOffset())); + .setOffset(toDelete.beforeOffset())); } DeleteRecordsRequestData data = new DeleteRecordsRequestData() diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/internals/DeleteRecordsHandlerTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/internals/DeleteRecordsHandlerTest.java index c39747f1fbafe..c773744031206 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/internals/DeleteRecordsHandlerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/internals/DeleteRecordsHandlerTest.java @@ -22,15 +22,24 @@ import java.util.List; import java.util.Map; import java.util.Set; +import java.util.stream.Collectors; + import org.apache.kafka.clients.admin.DeletedRecords; import org.apache.kafka.clients.admin.RecordsToDelete; +import org.apache.kafka.clients.admin.internals.AdminApiLookupStrategy.LookupResult; import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.message.DeleteRecordsRequestData; import org.apache.kafka.common.message.DeleteRecordsResponseData; +import org.apache.kafka.common.message.MetadataResponseData; +import org.apache.kafka.common.message.MetadataResponseData.MetadataResponsePartition; +import org.apache.kafka.common.message.MetadataResponseData.MetadataResponseTopic; +import org.apache.kafka.common.protocol.ApiKeys; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.requests.DeleteRecordsRequest; import org.apache.kafka.common.requests.DeleteRecordsResponse; +import org.apache.kafka.common.requests.MetadataRequest; +import org.apache.kafka.common.requests.MetadataResponse; import org.apache.kafka.common.utils.LogContext; import org.junit.jupiter.api.Test; @@ -41,6 +50,7 @@ import static org.apache.kafka.common.utils.Utils.mkSet; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertTrue; public class DeleteRecordsHandlerTest { @@ -50,7 +60,8 @@ public class DeleteRecordsHandlerTest { private final TopicPartition t0p1 = new TopicPartition("t0", 1); private final TopicPartition t0p2 = new TopicPartition("t0", 2); private final TopicPartition t0p3 = new TopicPartition("t0", 3); - private final Node node = new Node(1, "host", 1234); + private final Node node1 = new Node(1, "host", 1234); + private final Node node2 = new Node(2, "host", 1235); private final Map recordsToDelete = new HashMap() { { put(t0p0, RecordsToDelete.beforeOffset(10L)); @@ -63,11 +74,11 @@ public class DeleteRecordsHandlerTest { @Test public void testBuildRequestSimple() { DeleteRecordsHandler handler = new DeleteRecordsHandler(recordsToDelete, logContext, timeout); - DeleteRecordsRequest request = handler.buildBatchedRequest(node.id(), mkSet(t0p0, t0p1)).build(); + DeleteRecordsRequest request = handler.buildBatchedRequest(node1.id(), mkSet(t0p0, t0p1)).build(); List topicPartitions = request.data().topics(); assertEquals(1, topicPartitions.size()); DeleteRecordsRequestData.DeleteRecordsTopic topic = topicPartitions.get(0); - assertEquals(4, topic.partitions().size()); + assertEquals(2, topic.partitions().size()); } @Test @@ -199,6 +210,54 @@ public void testHandleResponseSanityCheck() { assertTrue(result.unmappedKeys.isEmpty()); } + // This is a more complicated test which ensures that DeleteRecords requests for multiple + // leader nodes are correctly divided up among the nodes based on leadership. + // node1 leads t0p0 and t0p2, while node2 leads t0p1 and t0p3. + @Test + public void testBuildRequestMultipleLeaders() { + MetadataResponseData metadataResponseData = new MetadataResponseData(); + MetadataResponseTopic topicMetadata = new MetadataResponseTopic(); + topicMetadata.setName("t0").setErrorCode(Errors.NONE.code()); + topicMetadata.partitions().add(new MetadataResponsePartition() + .setPartitionIndex(0).setLeaderId(node1.id()).setErrorCode(Errors.NONE.code())); + topicMetadata.partitions().add(new MetadataResponsePartition() + .setPartitionIndex(1).setLeaderId(node2.id()).setErrorCode(Errors.NONE.code())); + topicMetadata.partitions().add(new MetadataResponsePartition() + .setPartitionIndex(2).setLeaderId(node1.id()).setErrorCode(Errors.NONE.code())); + topicMetadata.partitions().add(new MetadataResponsePartition() + .setPartitionIndex(3).setLeaderId(node2.id()).setErrorCode(Errors.NONE.code())); + metadataResponseData.topics().add(topicMetadata); + MetadataResponse metadataResponse = new MetadataResponse(metadataResponseData, ApiKeys.METADATA.latestVersion()); + + DeleteRecordsHandler handler = new DeleteRecordsHandler(recordsToDelete, logContext, timeout); + AdminApiLookupStrategy strategy = handler.lookupStrategy(); + assertInstanceOf(PartitionLeaderStrategy.class, strategy); + PartitionLeaderStrategy specificStrategy = (PartitionLeaderStrategy) strategy; + MetadataRequest request = specificStrategy.buildRequest(mkSet(t0p0, t0p1, t0p2, t0p3)).build(); + assertEquals(mkSet("t0"), new HashSet<>(request.topics())); + + Set tpSet = mkSet(t0p0, t0p1, t0p2, t0p3); + LookupResult lookupResult = strategy.handleResponse(tpSet, metadataResponse); + assertEquals(emptyMap(), lookupResult.failedKeys); + assertEquals(tpSet, lookupResult.mappedKeys.keySet()); + + Map> flippedMapping = new HashMap<>(); + lookupResult.mappedKeys.forEach((tp, node) -> flippedMapping.computeIfAbsent(node, key -> new HashSet<>()).add(tp)); + + DeleteRecordsRequest deleteRequest = handler.buildBatchedRequest(node1.id(), flippedMapping.get(1)).build(); + assertEquals(2, deleteRequest.data().topics().get(0).partitions().size()); + assertEquals(mkSet(t0p0, t0p2), + new HashSet<>(deleteRequest.data().topics().get(0).partitions().stream() + .map(drp -> new TopicPartition("t0", drp.partitionIndex())) + .collect(Collectors.toList()))); + deleteRequest = handler.buildBatchedRequest(node2.id(), flippedMapping.get(2)).build(); + assertEquals(2, deleteRequest.data().topics().get(0).partitions().size()); + assertEquals(mkSet(t0p1, t0p3), + new HashSet<>(deleteRequest.data().topics().get(0).partitions().stream() + .map(drp -> new TopicPartition("t0", drp.partitionIndex())) + .collect(Collectors.toList()))); + } + private DeleteRecordsResponse createResponse(Map errorsByPartition) { return createResponse(errorsByPartition, recordsToDelete.keySet()); } @@ -227,7 +286,7 @@ private DeleteRecordsResponse createResponse( private AdminApiHandler.ApiResult handleResponse(DeleteRecordsResponse response) { DeleteRecordsHandler handler = new DeleteRecordsHandler(recordsToDelete, logContext, timeout); - return handler.handleResponse(node, recordsToDelete.keySet(), response); + return handler.handleResponse(node1, recordsToDelete.keySet(), response); } private void assertResult( From ecb789fc63c8500b3478774bfa5654a5452248e2 Mon Sep 17 00:00:00 2001 From: Andrew Schofield Date: Thu, 7 Mar 2024 10:09:48 +0000 Subject: [PATCH 2/2] Review comments --- .../internals/DeleteRecordsHandlerTest.java | 22 +++++++++---------- 1 file changed, 11 insertions(+), 11 deletions(-) diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/internals/DeleteRecordsHandlerTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/internals/DeleteRecordsHandlerTest.java index c773744031206..58492696c4da8 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/internals/DeleteRecordsHandlerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/internals/DeleteRecordsHandlerTest.java @@ -75,9 +75,9 @@ public class DeleteRecordsHandlerTest { public void testBuildRequestSimple() { DeleteRecordsHandler handler = new DeleteRecordsHandler(recordsToDelete, logContext, timeout); DeleteRecordsRequest request = handler.buildBatchedRequest(node1.id(), mkSet(t0p0, t0p1)).build(); - List topicPartitions = request.data().topics(); - assertEquals(1, topicPartitions.size()); - DeleteRecordsRequestData.DeleteRecordsTopic topic = topicPartitions.get(0); + List topics = request.data().topics(); + assertEquals(1, topics.size()); + DeleteRecordsRequestData.DeleteRecordsTopic topic = topics.get(0); assertEquals(2, topic.partitions().size()); } @@ -241,21 +241,21 @@ public void testBuildRequestMultipleLeaders() { assertEquals(emptyMap(), lookupResult.failedKeys); assertEquals(tpSet, lookupResult.mappedKeys.keySet()); - Map> flippedMapping = new HashMap<>(); - lookupResult.mappedKeys.forEach((tp, node) -> flippedMapping.computeIfAbsent(node, key -> new HashSet<>()).add(tp)); + Map> partitionsPerBroker = new HashMap<>(); + lookupResult.mappedKeys.forEach((tp, node) -> partitionsPerBroker.computeIfAbsent(node, key -> new HashSet<>()).add(tp)); - DeleteRecordsRequest deleteRequest = handler.buildBatchedRequest(node1.id(), flippedMapping.get(1)).build(); + DeleteRecordsRequest deleteRequest = handler.buildBatchedRequest(node1.id(), partitionsPerBroker.get(node1.id())).build(); assertEquals(2, deleteRequest.data().topics().get(0).partitions().size()); assertEquals(mkSet(t0p0, t0p2), - new HashSet<>(deleteRequest.data().topics().get(0).partitions().stream() + deleteRequest.data().topics().get(0).partitions().stream() .map(drp -> new TopicPartition("t0", drp.partitionIndex())) - .collect(Collectors.toList()))); - deleteRequest = handler.buildBatchedRequest(node2.id(), flippedMapping.get(2)).build(); + .collect(Collectors.toSet())); + deleteRequest = handler.buildBatchedRequest(node2.id(), partitionsPerBroker.get(node2.id())).build(); assertEquals(2, deleteRequest.data().topics().get(0).partitions().size()); assertEquals(mkSet(t0p1, t0p3), - new HashSet<>(deleteRequest.data().topics().get(0).partitions().stream() + deleteRequest.data().topics().get(0).partitions().stream() .map(drp -> new TopicPartition("t0", drp.partitionIndex())) - .collect(Collectors.toList()))); + .collect(Collectors.toSet())); } private DeleteRecordsResponse createResponse(Map errorsByPartition) {