Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 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 @@ -89,6 +89,11 @@
import org.apache.kafka.common.message.DeleteAclsResponseData.DeleteAclsFilterResult;
import org.apache.kafka.common.message.DeleteAclsResponseData.DeleteAclsMatchingAcl;
import org.apache.kafka.common.message.DeleteGroupsRequestData;
import org.apache.kafka.common.message.DeleteRecordsRequestData;
import org.apache.kafka.common.message.DeleteRecordsRequestData.DeleteRecordsPartition;
import org.apache.kafka.common.message.DeleteRecordsRequestData.DeleteRecordsTopic;
import org.apache.kafka.common.message.DeleteRecordsResponseData;
import org.apache.kafka.common.message.DeleteRecordsResponseData.DeleteRecordsTopicResult;
import org.apache.kafka.common.message.DeleteTopicsRequestData;
import org.apache.kafka.common.message.DeleteTopicsResponseData.DeletableTopicResult;
import org.apache.kafka.common.message.DescribeGroupsRequestData;
Expand Down Expand Up @@ -2469,19 +2474,30 @@ void handleResponse(AbstractResponse abstractResponse) {
Cluster cluster = response.cluster();

// Group topic partitions by leader
Map<Node, Map<TopicPartition, Long>> leaders = new HashMap<>();
Map<Node, Map<String, DeleteRecordsTopic>> leaders = new HashMap<>();
for (Map.Entry<TopicPartition, RecordsToDelete> entry: recordsToDelete.entrySet()) {
KafkaFutureImpl<DeletedRecords> future = futures.get(entry.getKey());
TopicPartition topicPartition = entry.getKey();
KafkaFutureImpl<DeletedRecords> future = futures.get(topicPartition);

// Fail partitions with topic errors
Errors topicError = errors.get(entry.getKey().topic());
if (errors.containsKey(entry.getKey().topic())) {
Errors topicError = errors.get(topicPartition.topic());
if (errors.containsKey(topicPartition.topic())) {
future.completeExceptionally(topicError.exception());
} else {
Node node = cluster.leaderFor(entry.getKey());
Node node = cluster.leaderFor(topicPartition);
if (node != null) {
leaders.computeIfAbsent(node, key -> new HashMap<>()).put(entry.getKey(),
entry.getValue().beforeOffset());
if (!leaders.containsKey(node))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I think we can computeIfAbsent() here. It will replace containsKey, put and get.

leaders.put(node, new HashMap<>());
Map<String, DeleteRecordsTopic> deletionsForLeader = leaders.get(node);
DeleteRecordsTopic deleteRecords = deletionsForLeader.get(topicPartition.topic());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Same here, I think computeIfAbsent would simplify the logic

if (deleteRecords == null) {
deleteRecords = new DeleteRecordsTopic()
.setName(topicPartition.topic());
deletionsForLeader.put(topicPartition.topic(), deleteRecords);
}
deleteRecords.partitions().add(new DeleteRecordsPartition()
.setPartitionIndex(topicPartition.partition())
.setOffset(entry.getValue().beforeOffset()));
} else {
future.completeExceptionally(Errors.LEADER_NOT_AVAILABLE.exception());
}
Expand All @@ -2490,36 +2506,45 @@ void handleResponse(AbstractResponse abstractResponse) {

final long deleteRecordsCallTimeMs = time.milliseconds();

for (final Map.Entry<Node, Map<TopicPartition, Long>> entry : leaders.entrySet()) {
final Map<TopicPartition, Long> partitionDeleteOffsets = entry.getValue();
for (final Map.Entry<Node, Map<String, DeleteRecordsTopic>> entry : leaders.entrySet()) {
final Map<String, DeleteRecordsTopic> partitionDeleteOffsets = entry.getValue();
final int brokerId = entry.getKey().id();

runnable.call(new Call("deleteRecords", deadline,
new ConstantNodeIdProvider(brokerId)) {

@Override
DeleteRecordsRequest.Builder createRequest(int timeoutMs) {
return new DeleteRecordsRequest.Builder(timeoutMs, partitionDeleteOffsets);
return new DeleteRecordsRequest.Builder(new DeleteRecordsRequestData()
.setTimeoutMs(timeoutMs)
.setTopics(new ArrayList<>(partitionDeleteOffsets.values())));
}

@Override
void handleResponse(AbstractResponse abstractResponse) {
DeleteRecordsResponse response = (DeleteRecordsResponse) abstractResponse;
for (Map.Entry<TopicPartition, DeleteRecordsResponse.PartitionResponse> result: response.responses().entrySet()) {

KafkaFutureImpl<DeletedRecords> future = futures.get(result.getKey());
if (result.getValue().error == Errors.NONE) {
future.complete(new DeletedRecords(result.getValue().lowWatermark));
} else {
future.completeExceptionally(result.getValue().error.exception());
for (DeleteRecordsTopicResult topicResult: response.data().topics()) {
for (DeleteRecordsResponseData.DeleteRecordsPartitionResult partitionResult : topicResult.partitions()) {
KafkaFutureImpl<DeletedRecords> future = futures.get(new TopicPartition(topicResult.name(), partitionResult.partitionIndex()));
if (partitionResult.errorCode() == Errors.NONE.code()) {
future.complete(new DeletedRecords(partitionResult.lowWatermark()));
} else {
future.completeExceptionally(Errors.forCode(partitionResult.errorCode()).exception());
}
}
}
}

@Override
void handleFailure(Throwable throwable) {
Stream<KafkaFutureImpl<DeletedRecords>> callFutures =
partitionDeleteOffsets.keySet().stream().map(futures::get);
partitionDeleteOffsets.values().stream().flatMap(
recordsToDelete1 -> {
Stream<TopicPartition> topicPartitionStream = recordsToDelete1.partitions().stream().map(partitionsToDelete ->

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Can we return that directly without creating a variable?

new TopicPartition(recordsToDelete1.name(), partitionsToDelete.partitionIndex()));
return topicPartitionStream;
}
).map(futures::get);
completeAllExceptionally(callFutures, throwable);
}
}, deleteRecordsCallTimeMs);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@
import org.apache.kafka.common.message.DeleteAclsResponseData;
import org.apache.kafka.common.message.DeleteGroupsRequestData;
import org.apache.kafka.common.message.DeleteGroupsResponseData;
import org.apache.kafka.common.message.DeleteRecordsRequestData;
import org.apache.kafka.common.message.DeleteRecordsResponseData;
import org.apache.kafka.common.message.DeleteTopicsRequestData;
import org.apache.kafka.common.message.DeleteTopicsResponseData;
import org.apache.kafka.common.message.DescribeAclsRequestData;
Expand Down Expand Up @@ -102,8 +104,6 @@
import org.apache.kafka.common.requests.AlterConfigsResponse;
import org.apache.kafka.common.requests.AlterReplicaLogDirsRequest;
import org.apache.kafka.common.requests.AlterReplicaLogDirsResponse;
import org.apache.kafka.common.requests.DeleteRecordsRequest;
import org.apache.kafka.common.requests.DeleteRecordsResponse;
import org.apache.kafka.common.requests.DescribeConfigsRequest;
import org.apache.kafka.common.requests.DescribeConfigsResponse;
import org.apache.kafka.common.requests.DescribeLogDirsRequest;
Expand Down Expand Up @@ -164,7 +164,7 @@ public Struct parseResponse(short version, ByteBuffer buffer) {
},
CREATE_TOPICS(19, "CreateTopics", CreateTopicsRequestData.SCHEMAS, CreateTopicsResponseData.SCHEMAS),
DELETE_TOPICS(20, "DeleteTopics", DeleteTopicsRequestData.SCHEMAS, DeleteTopicsResponseData.SCHEMAS),
DELETE_RECORDS(21, "DeleteRecords", DeleteRecordsRequest.schemaVersions(), DeleteRecordsResponse.schemaVersions()),
DELETE_RECORDS(21, "DeleteRecords", DeleteRecordsRequestData.SCHEMAS, DeleteRecordsResponseData.SCHEMAS),
INIT_PRODUCER_ID(22, "InitProducerId", InitProducerIdRequestData.SCHEMAS, InitProducerIdResponseData.SCHEMAS),
OFFSET_FOR_LEADER_EPOCH(23, "OffsetForLeaderEpoch", false, OffsetsForLeaderEpochRequest.schemaVersions(),
OffsetsForLeaderEpochResponse.schemaVersions()),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -123,7 +123,7 @@ public static AbstractResponse parseResponse(ApiKeys apiKey, Struct struct, shor
case DELETE_TOPICS:
return new DeleteTopicsResponse(struct, version);
case DELETE_RECORDS:
return new DeleteRecordsResponse(struct);
return new DeleteRecordsResponse(struct, version);
case INIT_PRODUCER_ID:
return new InitProducerIdResponse(struct, version);
case OFFSET_FOR_LEADER_EPOCH:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,153 +17,75 @@

package org.apache.kafka.common.requests;

import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.message.DeleteRecordsRequestData;
import org.apache.kafka.common.message.DeleteRecordsRequestData.DeleteRecordsTopic;
import org.apache.kafka.common.message.DeleteRecordsResponseData;
import org.apache.kafka.common.message.DeleteRecordsResponseData.DeleteRecordsTopicResult;
import org.apache.kafka.common.protocol.ApiKeys;
import org.apache.kafka.common.protocol.Errors;
import org.apache.kafka.common.protocol.types.ArrayOf;
import org.apache.kafka.common.protocol.types.Field;
import org.apache.kafka.common.protocol.types.Schema;
import org.apache.kafka.common.protocol.types.Struct;
import org.apache.kafka.common.utils.CollectionUtils;

import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;

import static org.apache.kafka.common.protocol.CommonFields.PARTITION_ID;
import static org.apache.kafka.common.protocol.CommonFields.TOPIC_NAME;
import static org.apache.kafka.common.protocol.types.Type.INT32;
import static org.apache.kafka.common.protocol.types.Type.INT64;

public class DeleteRecordsRequest extends AbstractRequest {

public static final long HIGH_WATERMARK = -1L;

// request level key names
private static final String TOPICS_KEY_NAME = "topics";
private static final String TIMEOUT_KEY_NAME = "timeout";

// topic level key names
private static final String PARTITIONS_KEY_NAME = "partitions";

// partition level key names
private static final String OFFSET_KEY_NAME = "offset";


private static final Schema DELETE_RECORDS_REQUEST_PARTITION_V0 = new Schema(
PARTITION_ID,
new Field(OFFSET_KEY_NAME, INT64, "The offset before which the messages will be deleted. -1 means high-watermark for the partition."));

private static final Schema DELETE_RECORDS_REQUEST_TOPIC_V0 = new Schema(
TOPIC_NAME,
new Field(PARTITIONS_KEY_NAME, new ArrayOf(DELETE_RECORDS_REQUEST_PARTITION_V0)));

private static final Schema DELETE_RECORDS_REQUEST_V0 = new Schema(
new Field(TOPICS_KEY_NAME, new ArrayOf(DELETE_RECORDS_REQUEST_TOPIC_V0)),
new Field(TIMEOUT_KEY_NAME, INT32, "The maximum time to await a response in ms."));

/**
* The version number is bumped to indicate that on quota violation brokers send out responses before throttling.
*/
private static final Schema DELETE_RECORDS_REQUEST_V1 = DELETE_RECORDS_REQUEST_V0;

public static Schema[] schemaVersions() {
return new Schema[]{DELETE_RECORDS_REQUEST_V0, DELETE_RECORDS_REQUEST_V1};
}

private final int timeout;
private final Map<TopicPartition, Long> partitionOffsets;
private final DeleteRecordsRequestData data;

public static class Builder extends AbstractRequest.Builder<DeleteRecordsRequest> {
private final int timeout;
private final Map<TopicPartition, Long> partitionOffsets;
private DeleteRecordsRequestData data;

public Builder(int timeout, Map<TopicPartition, Long> partitionOffsets) {
public Builder(DeleteRecordsRequestData data) {
super(ApiKeys.DELETE_RECORDS);
this.timeout = timeout;
this.partitionOffsets = partitionOffsets;
this.data = data;
}

@Override
public DeleteRecordsRequest build(short version) {
return new DeleteRecordsRequest(timeout, partitionOffsets, version);
return new DeleteRecordsRequest(data, version);
}

@Override
public String toString() {
StringBuilder builder = new StringBuilder();
builder.append("(type=DeleteRecordsRequest")
.append(", timeout=").append(timeout)
.append(", partitionOffsets=(").append(partitionOffsets)
.append("))");
return builder.toString();
return data.toString();
}
}


public DeleteRecordsRequest(Struct struct, short version) {
private DeleteRecordsRequest(DeleteRecordsRequestData data, short version) {
super(ApiKeys.DELETE_RECORDS, version);
partitionOffsets = new HashMap<>();
for (Object topicStructObj : struct.getArray(TOPICS_KEY_NAME)) {
Struct topicStruct = (Struct) topicStructObj;
String topic = topicStruct.get(TOPIC_NAME);
for (Object partitionStructObj : topicStruct.getArray(PARTITIONS_KEY_NAME)) {
Struct partitionStruct = (Struct) partitionStructObj;
int partition = partitionStruct.get(PARTITION_ID);
long offset = partitionStruct.getLong(OFFSET_KEY_NAME);
partitionOffsets.put(new TopicPartition(topic, partition), offset);
}
}
timeout = struct.getInt(TIMEOUT_KEY_NAME);
this.data = data;
}

public DeleteRecordsRequest(int timeout, Map<TopicPartition, Long> partitionOffsets, short version) {
public DeleteRecordsRequest(Struct struct, short version) {
super(ApiKeys.DELETE_RECORDS, version);
this.timeout = timeout;
this.partitionOffsets = partitionOffsets;
this.data = new DeleteRecordsRequestData(struct, version);
}

@Override
protected Struct toStruct() {
Struct struct = new Struct(ApiKeys.DELETE_RECORDS.requestSchema(version()));
Map<String, Map<Integer, Long>> offsetsByTopic = CollectionUtils.groupPartitionDataByTopic(partitionOffsets);
struct.set(TIMEOUT_KEY_NAME, timeout);
List<Struct> topicStructArray = new ArrayList<>();
for (Map.Entry<String, Map<Integer, Long>> offsetsByTopicEntry : offsetsByTopic.entrySet()) {
Struct topicStruct = struct.instance(TOPICS_KEY_NAME);
topicStruct.set(TOPIC_NAME, offsetsByTopicEntry.getKey());
List<Struct> partitionStructArray = new ArrayList<>();
for (Map.Entry<Integer, Long> offsetsByPartitionEntry : offsetsByTopicEntry.getValue().entrySet()) {
Struct partitionStruct = topicStruct.instance(PARTITIONS_KEY_NAME);
partitionStruct.set(PARTITION_ID, offsetsByPartitionEntry.getKey());
partitionStruct.set(OFFSET_KEY_NAME, offsetsByPartitionEntry.getValue());
partitionStructArray.add(partitionStruct);
}
topicStruct.set(PARTITIONS_KEY_NAME, partitionStructArray.toArray());
topicStructArray.add(topicStruct);
}
struct.set(TOPICS_KEY_NAME, topicStructArray.toArray());
return struct;
return data.toStruct(version());
}

public DeleteRecordsRequestData data() {
return data;
}

@Override
public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
Map<TopicPartition, DeleteRecordsResponse.PartitionResponse> responseMap = new HashMap<>();

for (Map.Entry<TopicPartition, Long> entry : partitionOffsets.entrySet()) {
responseMap.put(entry.getKey(), new DeleteRecordsResponse.PartitionResponse(DeleteRecordsResponse.INVALID_LOW_WATERMARK, Errors.forException(e)));
DeleteRecordsResponseData result = new DeleteRecordsResponseData().setThrottleTimeMs(throttleTimeMs);
short errorCode = Errors.forException(e).code();
for (DeleteRecordsTopic topic : data.topics()) {
DeleteRecordsTopicResult topicResult = new DeleteRecordsTopicResult().setName(topic.name());
result.topics().add(topicResult);
for (DeleteRecordsRequestData.DeleteRecordsPartition partition : topic.partitions()) {
topicResult.partitions().add(new DeleteRecordsResponseData.DeleteRecordsPartitionResult()
.setPartitionIndex(partition.partitionIndex())
.setErrorCode(errorCode)
.setLowWatermark(DeleteRecordsResponse.INVALID_LOW_WATERMARK));
}
}

return new DeleteRecordsResponse(throttleTimeMs, responseMap);
}

public int timeout() {
return timeout;
}

public Map<TopicPartition, Long> partitionOffsets() {
return partitionOffsets;
return new DeleteRecordsResponse(result);
}

public static DeleteRecordsRequest parse(ByteBuffer buffer, short version) {
Expand Down
Loading