Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -79,15 +79,15 @@ public static SimpleAdminApiFuture<TopicPartition, DeletedRecords> newFuture(
@Override
public DeleteRecordsRequest.Builder buildBatchedRequest(int brokerId, Set<TopicPartition> keys) {
Map<String, DeleteRecordsRequestData.DeleteRecordsTopic> deletionsForTopic = new HashMap<>();
for (Map.Entry<TopicPartition, RecordsToDelete> entry: recordsToDelete.entrySet()) {
TopicPartition topicPartition = entry.getKey();
for (TopicPartition topicPartition : keys) {
Comment thread
AndrewJSchofield marked this conversation as resolved.
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()));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would it be unduly paranoid to check for null before accessing toDelete? The call stack up to this point is pretty twisty. I'm assuming that keys shouldn't contain any TopicPartitions that aren't keys in recordsToDelete, but... 🤷‍♂️

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, in my opinion it would be unduly paranoid. I would end up writing conditional logic for a situation which will never occur, and it would differ from other code in this area of AdminClient which also uses a map which is accessed in this way confidently expecting the entry to be present.

}

DeleteRecordsRequestData data = new DeleteRecordsRequestData()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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 {
Expand All @@ -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<TopicPartition, RecordsToDelete> recordsToDelete = new HashMap<TopicPartition, RecordsToDelete>() {
{
put(t0p0, RecordsToDelete.beforeOffset(10L));
Expand All @@ -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();
List<DeleteRecordsRequestData.DeleteRecordsTopic> topicPartitions = request.data().topics();
assertEquals(1, topicPartitions.size());
DeleteRecordsRequestData.DeleteRecordsTopic topic = topicPartitions.get(0);
assertEquals(4, topic.partitions().size());
DeleteRecordsRequest request = handler.buildBatchedRequest(node1.id(), mkSet(t0p0, t0p1)).build();
List<DeleteRecordsRequestData.DeleteRecordsTopic> topics = request.data().topics();
assertEquals(1, topics.size());
DeleteRecordsRequestData.DeleteRecordsTopic topic = topics.get(0);
assertEquals(2, topic.partitions().size());
}

@Test
Expand Down Expand Up @@ -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<TopicPartition> 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<TopicPartition> tpSet = mkSet(t0p0, t0p1, t0p2, t0p3);
LookupResult<TopicPartition> lookupResult = strategy.handleResponse(tpSet, metadataResponse);
assertEquals(emptyMap(), lookupResult.failedKeys);
assertEquals(tpSet, lookupResult.mappedKeys.keySet());

Map<Integer, Set<TopicPartition>> partitionsPerBroker = new HashMap<>();
lookupResult.mappedKeys.forEach((tp, node) -> partitionsPerBroker.computeIfAbsent(node, key -> new HashSet<>()).add(tp));

DeleteRecordsRequest deleteRequest = handler.buildBatchedRequest(node1.id(), partitionsPerBroker.get(node1.id())).build();
assertEquals(2, deleteRequest.data().topics().get(0).partitions().size());
assertEquals(mkSet(t0p0, t0p2),
deleteRequest.data().topics().get(0).partitions().stream()
.map(drp -> new TopicPartition("t0", drp.partitionIndex()))
.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),
deleteRequest.data().topics().get(0).partitions().stream()
.map(drp -> new TopicPartition("t0", drp.partitionIndex()))
.collect(Collectors.toSet()));
}

private DeleteRecordsResponse createResponse(Map<TopicPartition, Short> errorsByPartition) {
return createResponse(errorsByPartition, recordsToDelete.keySet());
}
Expand Down Expand Up @@ -227,7 +286,7 @@ private DeleteRecordsResponse createResponse(
private AdminApiHandler.ApiResult<TopicPartition, DeletedRecords> 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(
Expand Down