diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml
index 1bb797bd4497c..d57c760d5bf41 100644
--- a/checkstyle/suppressions.xml
+++ b/checkstyle/suppressions.xml
@@ -62,7 +62,7 @@
files="(ConsumerCoordinator|Fetcher|Sender|KafkaProducer|BufferPool|ConfigDef|RecordAccumulator|KerberosLogin|AbstractRequest|AbstractResponse|Selector|SslFactory|SslTransportLayer|SaslClientAuthenticator|SaslClientCallbackHandler|SaslServerAuthenticator|AbstractCoordinator|TransactionManager).java"/>
+ files="(AbstractRequest|AbstractResponse|KerberosLogin|WorkerSinkTaskTest|TransactionManagerTest|SenderTest|KafkaAdminClient|ConsumerCoordinatorTest|KafkaAdminClientTest).java"/>
diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
index 7d2a4d09914df..3dbd3007999fd 100644
--- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
+++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
@@ -69,6 +69,11 @@
import org.apache.kafka.common.internals.KafkaFutureImpl;
import org.apache.kafka.common.message.AlterPartitionReassignmentsRequestData;
import org.apache.kafka.common.message.AlterPartitionReassignmentsRequestData.ReassignableTopic;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData.AlterReplicaLogDir;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData.AlterReplicaLogDirTopic;
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData.AlterReplicaLogDirPartitionResult;
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData.AlterReplicaLogDirTopicResult;
import org.apache.kafka.common.message.CreateAclsRequestData;
import org.apache.kafka.common.message.CreateAclsRequestData.AclCreation;
import org.apache.kafka.common.message.CreateAclsResponseData.AclCreationResult;
@@ -234,6 +239,7 @@
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
+import java.util.function.Function;
import java.util.function.Predicate;
import java.util.function.Supplier;
import java.util.stream.Collectors;
@@ -1402,6 +1408,17 @@ int numPendingCalls() {
return runnable.pendingCalls.size();
}
+ /**
+ * Fail futures in the given stream which are not done.
+ * Used when a response handler expected a result for some entity but no result was present.
+ */
+ private static void completeUnrealizedFutures(
+ Stream>> futures,
+ Function messageFormatter) {
+ futures.filter(entry -> !entry.getValue().isDone()).forEach(entry ->
+ entry.getValue().completeExceptionally(new ApiException(messageFormatter.apply(entry.getKey()))));
+ }
+
@Override
public CreateTopicsResult createTopics(final Collection newTopics,
final CreateTopicsOptions options) {
@@ -1479,13 +1496,8 @@ public void handleResponse(AbstractResponse abstractResponse) {
}
}
// The server should send back a response for every topic. But do a sanity check anyway.
- for (Map.Entry> entry : topicFutures.entrySet()) {
- KafkaFutureImpl future = entry.getValue();
- if (!future.isDone()) {
- future.completeExceptionally(new ApiException("The server response did not " +
- "contain a reference to node " + entry.getKey()));
- }
- }
+ completeUnrealizedFutures(topicFutures.entrySet().stream(),
+ topic -> "The controller response did not contain a result for topic " + topic);
}
@Override
@@ -1552,13 +1564,8 @@ void handleResponse(AbstractResponse abstractResponse) {
}
}
// The server should send back a response for every topic. But do a sanity check anyway.
- for (Map.Entry> entry : topicFutures.entrySet()) {
- KafkaFutureImpl future = entry.getValue();
- if (!future.isDone()) {
- future.completeExceptionally(new ApiException("The server response did not " +
- "contain a reference to node " + entry.getKey()));
- }
- }
+ completeUnrealizedFutures(topicFutures.entrySet().stream(),
+ topic -> "The controller response did not contain a result for topic " + topic);
}
@Override
@@ -2207,21 +2214,31 @@ public AlterReplicaLogDirsResult alterReplicaLogDirs(Map());
- Map> replicaAssignmentByBroker = new HashMap<>();
+ Map replicaAssignmentByBroker = new HashMap<>();
for (Map.Entry entry: replicaAssignment.entrySet()) {
TopicPartitionReplica replica = entry.getKey();
String logDir = entry.getValue();
int brokerId = replica.brokerId();
- TopicPartition topicPartition = new TopicPartition(replica.topic(), replica.partition());
- if (!replicaAssignmentByBroker.containsKey(brokerId))
- replicaAssignmentByBroker.put(brokerId, new HashMap<>());
- replicaAssignmentByBroker.get(brokerId).put(topicPartition, logDir);
+ AlterReplicaLogDirsRequestData value = replicaAssignmentByBroker.computeIfAbsent(brokerId,
+ key -> new AlterReplicaLogDirsRequestData());
+ AlterReplicaLogDir alterReplicaLogDir = value.dirs().find(logDir);
+ if (alterReplicaLogDir == null) {
+ alterReplicaLogDir = new AlterReplicaLogDir();
+ alterReplicaLogDir.setPath(logDir);
+ value.dirs().add(alterReplicaLogDir);
+ }
+ AlterReplicaLogDirTopic alterReplicaLogDirTopic = alterReplicaLogDir.topics().find(replica.topic());
+ if (alterReplicaLogDirTopic == null) {
+ alterReplicaLogDirTopic = new AlterReplicaLogDirTopic().setName(replica.topic());
+ alterReplicaLogDir.topics().add(alterReplicaLogDirTopic);
+ }
+ alterReplicaLogDirTopic.partitions().add(replica.partition());
}
final long now = time.milliseconds();
- for (Map.Entry> entry: replicaAssignmentByBroker.entrySet()) {
+ for (Map.Entry entry: replicaAssignmentByBroker.entrySet()) {
final int brokerId = entry.getKey();
- final Map assignment = entry.getValue();
+ final AlterReplicaLogDirsRequestData assignment = entry.getValue();
runnable.call(new Call("alterReplicaLogDirs", calcDeadlineMs(now, options.timeoutMs()),
new ConstantNodeIdProvider(brokerId)) {
@@ -2234,20 +2251,27 @@ public AlterReplicaLogDirsRequest.Builder createRequest(int timeoutMs) {
@Override
public void handleResponse(AbstractResponse abstractResponse) {
AlterReplicaLogDirsResponse response = (AlterReplicaLogDirsResponse) abstractResponse;
- for (Map.Entry responseEntry: response.responses().entrySet()) {
- TopicPartition tp = responseEntry.getKey();
- Errors error = responseEntry.getValue();
- TopicPartitionReplica replica = new TopicPartitionReplica(tp.topic(), tp.partition(), brokerId);
- KafkaFutureImpl future = futures.get(replica);
- if (future == null) {
- handleFailure(new IllegalStateException(
- "The partition " + tp + " in the response from broker " + brokerId + " is not in the request"));
- } else if (error == Errors.NONE) {
- future.complete(null);
- } else {
- future.completeExceptionally(error.exception());
+ for (AlterReplicaLogDirTopicResult topicResult: response.data().results()) {
+ for (AlterReplicaLogDirPartitionResult partitionResult: topicResult.partitions()) {
+ TopicPartitionReplica replica = new TopicPartitionReplica(
+ topicResult.topicName(), partitionResult.partitionIndex(), brokerId);
+ KafkaFutureImpl future = futures.get(replica);
+ if (future == null) {
+ log.warn("The partition {} in the response from broker {} is not in the request",
+ new TopicPartition(topicResult.topicName(), partitionResult.partitionIndex()),
+ brokerId);
+ } else if (partitionResult.errorCode() == Errors.NONE.code()) {
+ future.complete(null);
+ } else {
+ future.completeExceptionally(Errors.forCode(partitionResult.errorCode()).exception());
+ }
}
}
+ // The server should send back a response for every replica. But do a sanity check anyway.
+ completeUnrealizedFutures(
+ futures.entrySet().stream().filter(entry -> entry.getKey().brokerId() == brokerId),
+ replica -> "The response from broker " + brokerId +
+ " did not contain a result for replica " + replica);
}
@Override
void handleFailure(Throwable throwable) {
diff --git a/clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java b/clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java
index ee78f68098795..4cb9119d325a0 100644
--- a/clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java
+++ b/clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java
@@ -20,6 +20,8 @@
import org.apache.kafka.common.message.AddPartitionsToTxnResponseData;
import org.apache.kafka.common.message.AlterPartitionReassignmentsRequestData;
import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData;
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData;
import org.apache.kafka.common.message.ApiMessageType;
import org.apache.kafka.common.message.AddOffsetsToTxnRequestData;
import org.apache.kafka.common.message.AddOffsetsToTxnResponseData;
@@ -110,8 +112,6 @@
import org.apache.kafka.common.protocol.types.Struct;
import org.apache.kafka.common.protocol.types.Type;
import org.apache.kafka.common.record.RecordBatch;
-import org.apache.kafka.common.requests.AlterReplicaLogDirsRequest;
-import org.apache.kafka.common.requests.AlterReplicaLogDirsResponse;
import org.apache.kafka.common.requests.DescribeConfigsRequest;
import org.apache.kafka.common.requests.DescribeConfigsResponse;
import org.apache.kafka.common.requests.FetchRequest;
@@ -188,8 +188,8 @@ public Struct parseResponse(short version, ByteBuffer buffer) {
DescribeConfigsResponse.schemaVersions()),
ALTER_CONFIGS(33, "AlterConfigs", AlterConfigsRequestData.SCHEMAS,
AlterConfigsResponseData.SCHEMAS),
- ALTER_REPLICA_LOG_DIRS(34, "AlterReplicaLogDirs", AlterReplicaLogDirsRequest.schemaVersions(),
- AlterReplicaLogDirsResponse.schemaVersions()),
+ ALTER_REPLICA_LOG_DIRS(34, "AlterReplicaLogDirs", AlterReplicaLogDirsRequestData.SCHEMAS,
+ AlterReplicaLogDirsResponseData.SCHEMAS),
DESCRIBE_LOG_DIRS(35, "DescribeLogDirs", DescribeLogDirsRequestData.SCHEMAS,
DescribeLogDirsResponseData.SCHEMAS),
SASL_AUTHENTICATE(36, "SaslAuthenticate", SaslAuthenticateRequestData.SCHEMAS,
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterReplicaLogDirsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterReplicaLogDirsRequest.java
index 888470f259626..8cde2c0a1bde2 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AlterReplicaLogDirsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterReplicaLogDirsRequest.java
@@ -17,141 +17,79 @@
package org.apache.kafka.common.requests;
-import org.apache.kafka.common.TopicPartition;
-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 java.util.stream.Collectors;
-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.STRING;
+import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData;
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData;
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData.AlterReplicaLogDirTopicResult;
+import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.Errors;
+import org.apache.kafka.common.protocol.types.Struct;
public class AlterReplicaLogDirsRequest extends AbstractRequest {
- // request level key names
- private static final String LOG_DIRS_KEY_NAME = "log_dirs";
-
- // log dir level key names
- private static final String LOG_DIR_KEY_NAME = "log_dir";
- private static final String TOPICS_KEY_NAME = "topics";
-
- // topic level key names
- private static final String PARTITIONS_KEY_NAME = "partitions";
-
- private static final Schema ALTER_REPLICA_LOG_DIRS_REQUEST_V0 = new Schema(
- new Field("log_dirs", new ArrayOf(new Schema(
- new Field("log_dir", STRING, "The absolute log directory path."),
- new Field("topics", new ArrayOf(new Schema(
- TOPIC_NAME,
- new Field("partitions", new ArrayOf(INT32), "List of partition ids of the topic."))))))));
-
- /**
- * The version number is bumped to indicate that on quota violation brokers send out responses before throttling.
- */
- private static final Schema ALTER_REPLICA_LOG_DIRS_REQUEST_V1 = ALTER_REPLICA_LOG_DIRS_REQUEST_V0;
-
- public static Schema[] schemaVersions() {
- return new Schema[]{ALTER_REPLICA_LOG_DIRS_REQUEST_V0, ALTER_REPLICA_LOG_DIRS_REQUEST_V1};
- }
-
- private final Map partitionDirs;
+ private final AlterReplicaLogDirsRequestData data;
public static class Builder extends AbstractRequest.Builder {
- private final Map partitionDirs;
+ private final AlterReplicaLogDirsRequestData data;
- public Builder(Map partitionDirs) {
+ public Builder(AlterReplicaLogDirsRequestData data) {
super(ApiKeys.ALTER_REPLICA_LOG_DIRS);
- this.partitionDirs = partitionDirs;
+ this.data = data;
}
@Override
public AlterReplicaLogDirsRequest build(short version) {
- return new AlterReplicaLogDirsRequest(partitionDirs, version);
+ return new AlterReplicaLogDirsRequest(data, version);
}
@Override
public String toString() {
- StringBuilder builder = new StringBuilder();
- builder.append("(type=AlterReplicaLogDirsRequest")
- .append(", partitionDirs=")
- .append(partitionDirs)
- .append(")");
- return builder.toString();
+ return data.toString();
}
}
public AlterReplicaLogDirsRequest(Struct struct, short version) {
super(ApiKeys.ALTER_REPLICA_LOG_DIRS, version);
- partitionDirs = new HashMap<>();
- for (Object logDirStructObj : struct.getArray(LOG_DIRS_KEY_NAME)) {
- Struct logDirStruct = (Struct) logDirStructObj;
- String logDir = logDirStruct.getString(LOG_DIR_KEY_NAME);
- for (Object topicStructObj : logDirStruct.getArray(TOPICS_KEY_NAME)) {
- Struct topicStruct = (Struct) topicStructObj;
- String topic = topicStruct.get(TOPIC_NAME);
- for (Object partitionObj : topicStruct.getArray(PARTITIONS_KEY_NAME)) {
- int partition = (Integer) partitionObj;
- partitionDirs.put(new TopicPartition(topic, partition), logDir);
- }
- }
- }
+ this.data = new AlterReplicaLogDirsRequestData(struct, version);
}
- public AlterReplicaLogDirsRequest(Map partitionDirs, short version) {
+ public AlterReplicaLogDirsRequest(AlterReplicaLogDirsRequestData data, short version) {
super(ApiKeys.ALTER_REPLICA_LOG_DIRS, version);
- this.partitionDirs = partitionDirs;
+ this.data = data;
}
@Override
protected Struct toStruct() {
- Map> dirPartitions = new HashMap<>();
- for (Map.Entry entry: partitionDirs.entrySet()) {
- if (!dirPartitions.containsKey(entry.getValue()))
- dirPartitions.put(entry.getValue(), new ArrayList<>());
- dirPartitions.get(entry.getValue()).add(entry.getKey());
- }
-
- Struct struct = new Struct(ApiKeys.ALTER_REPLICA_LOG_DIRS.requestSchema(version()));
- List logDirStructArray = new ArrayList<>();
- for (Map.Entry> logDirEntry: dirPartitions.entrySet()) {
- Struct logDirStruct = struct.instance(LOG_DIRS_KEY_NAME);
- logDirStruct.set(LOG_DIR_KEY_NAME, logDirEntry.getKey());
-
- List topicStructArray = new ArrayList<>();
- for (Map.Entry> topicEntry: CollectionUtils.groupPartitionsByTopic(logDirEntry.getValue()).entrySet()) {
- Struct topicStruct = logDirStruct.instance(TOPICS_KEY_NAME);
- topicStruct.set(TOPIC_NAME, topicEntry.getKey());
- topicStruct.set(PARTITIONS_KEY_NAME, topicEntry.getValue().toArray());
- topicStructArray.add(topicStruct);
- }
- logDirStruct.set(TOPICS_KEY_NAME, topicStructArray.toArray());
- logDirStructArray.add(logDirStruct);
- }
- struct.set(LOG_DIRS_KEY_NAME, logDirStructArray.toArray());
- return struct;
+ return data.toStruct(version());
}
@Override
- public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
- Map responseMap = new HashMap<>();
- for (Map.Entry entry : partitionDirs.entrySet()) {
- responseMap.put(entry.getKey(), Errors.forException(e));
- }
- return new AlterReplicaLogDirsResponse(throttleTimeMs, responseMap);
+ public AlterReplicaLogDirsResponse getErrorResponse(int throttleTimeMs, Throwable e) {
+ AlterReplicaLogDirsResponseData data = new AlterReplicaLogDirsResponseData();
+ data.setResults(this.data.dirs().stream().flatMap(alterDir ->
+ alterDir.topics().stream().map(topic ->
+ new AlterReplicaLogDirTopicResult()
+ .setTopicName(topic.name())
+ .setPartitions(topic.partitions().stream().map(partitionId ->
+ new AlterReplicaLogDirsResponseData.AlterReplicaLogDirPartitionResult()
+ .setErrorCode(Errors.forException(e).code())
+ .setPartitionIndex(partitionId)).collect(Collectors.toList())))).collect(Collectors.toList()));
+ return new AlterReplicaLogDirsResponse(data.setThrottleTimeMs(throttleTimeMs));
}
public Map partitionDirs() {
- return partitionDirs;
+ Map result = new HashMap<>();
+ data.dirs().forEach(alterDir ->
+ alterDir.topics().forEach(topic ->
+ topic.partitions().forEach(partition ->
+ result.put(new TopicPartition(topic.name(), partition.intValue()), alterDir.path())))
+ );
+ return result;
}
public static AlterReplicaLogDirsRequest parse(ByteBuffer buffer, short version) {
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterReplicaLogDirsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterReplicaLogDirsResponse.java
index 2bbdce2122bc2..98d95c34dcd3f 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AlterReplicaLogDirsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterReplicaLogDirsResponse.java
@@ -17,122 +17,60 @@
package org.apache.kafka.common.requests;
-import org.apache.kafka.common.TopicPartition;
-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.ERROR_CODE;
-import static org.apache.kafka.common.protocol.CommonFields.PARTITION_ID;
-import static org.apache.kafka.common.protocol.CommonFields.THROTTLE_TIME_MS;
-import static org.apache.kafka.common.protocol.CommonFields.TOPIC_NAME;
-
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData;
+import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.Errors;
+import org.apache.kafka.common.protocol.types.Struct;
+/**
+ * Possible error codes:
+ *
+ * {@link Errors#LOG_DIR_NOT_FOUND}
+ * {@link Errors#KAFKA_STORAGE_ERROR}
+ * {@link Errors#REPLICA_NOT_AVAILABLE}
+ * {@link Errors#UNKNOWN_SERVER_ERROR}
+ */
public class AlterReplicaLogDirsResponse extends AbstractResponse {
- // request level key names
- private static final String TOPICS_KEY_NAME = "topics";
-
- // topic level key names
- private static final String PARTITIONS_KEY_NAME = "partitions";
-
- private static final Schema ALTER_REPLICA_LOG_DIRS_RESPONSE_V0 = new Schema(
- THROTTLE_TIME_MS,
- new Field(TOPICS_KEY_NAME, new ArrayOf(new Schema(
- TOPIC_NAME,
- new Field(PARTITIONS_KEY_NAME, new ArrayOf(new Schema(
- PARTITION_ID,
- ERROR_CODE)))))));
+ private final AlterReplicaLogDirsResponseData data;
- /**
- * The version number is bumped to indicate that on quota violation brokers send out responses before throttling.
- */
- private static final Schema ALTER_REPLICA_LOG_DIRS_RESPONSE_V1 = ALTER_REPLICA_LOG_DIRS_RESPONSE_V0;
-
- public static Schema[] schemaVersions() {
- return new Schema[]{ALTER_REPLICA_LOG_DIRS_RESPONSE_V0, ALTER_REPLICA_LOG_DIRS_RESPONSE_V1};
+ public AlterReplicaLogDirsResponse(Struct struct) {
+ this(struct, ApiKeys.ALTER_REPLICA_LOG_DIRS.latestVersion());
}
- /**
- * Possible error code:
- *
- * LOG_DIR_NOT_FOUND (57)
- * KAFKA_STORAGE_ERROR (56)
- * REPLICA_NOT_AVAILABLE (9)
- * UNKNOWN (-1)
- */
- private final Map responses;
- private final int throttleTimeMs;
+ public AlterReplicaLogDirsResponse(Struct struct, short version) {
+ this.data = new AlterReplicaLogDirsResponseData(struct, version);
+ }
- public AlterReplicaLogDirsResponse(Struct struct) {
- throttleTimeMs = struct.get(THROTTLE_TIME_MS);
- responses = 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);
- Errors error = Errors.forCode(partitionStruct.get(ERROR_CODE));
- responses.put(new TopicPartition(topic, partition), error);
- }
- }
+ public AlterReplicaLogDirsResponse(AlterReplicaLogDirsResponseData data) {
+ this.data = data;
}
- /**
- * Constructor for version 0.
- */
- public AlterReplicaLogDirsResponse(int throttleTimeMs, Map responses) {
- this.throttleTimeMs = throttleTimeMs;
- this.responses = responses;
+ public AlterReplicaLogDirsResponseData data() {
+ return data;
}
@Override
protected Struct toStruct(short version) {
- Struct struct = new Struct(ApiKeys.ALTER_REPLICA_LOG_DIRS.responseSchema(version));
- struct.set(THROTTLE_TIME_MS, throttleTimeMs);
- Map> responsesByTopic = CollectionUtils.groupPartitionDataByTopic(responses);
- List topicStructArray = new ArrayList<>();
- for (Map.Entry> responsesByTopicEntry : responsesByTopic.entrySet()) {
- Struct topicStruct = struct.instance(TOPICS_KEY_NAME);
- topicStruct.set(TOPIC_NAME, responsesByTopicEntry.getKey());
- List partitionStructArray = new ArrayList<>();
- for (Map.Entry responsesByPartitionEntry : responsesByTopicEntry.getValue().entrySet()) {
- Struct partitionStruct = topicStruct.instance(PARTITIONS_KEY_NAME);
- Errors response = responsesByPartitionEntry.getValue();
- partitionStruct.set(PARTITION_ID, responsesByPartitionEntry.getKey());
- partitionStruct.set(ERROR_CODE, response.code());
- 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);
}
@Override
public int throttleTimeMs() {
- return throttleTimeMs;
- }
-
- public Map responses() {
- return this.responses;
+ return data.throttleTimeMs();
}
@Override
public Map errorCounts() {
- return errorCounts(responses.values());
+ Map errorCounts = new HashMap<>();
+ data.results().forEach(topicResult ->
+ topicResult.partitions().forEach(partitionResult ->
+ updateErrorCounts(errorCounts, Errors.forCode(partitionResult.errorCode()))));
+ return errorCounts;
}
public static AlterReplicaLogDirsResponse parse(ByteBuffer buffer, short version) {
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 2510ea727d43d..333e3f6fe7bf2 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
@@ -32,6 +32,7 @@
import org.apache.kafka.common.Node;
import org.apache.kafka.common.PartitionInfo;
import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.TopicPartitionReplica;
import org.apache.kafka.common.acl.AccessControlEntry;
import org.apache.kafka.common.acl.AccessControlEntryFilter;
import org.apache.kafka.common.acl.AclBinding;
@@ -40,6 +41,7 @@
import org.apache.kafka.common.acl.AclPermissionType;
import org.apache.kafka.common.config.ConfigException;
import org.apache.kafka.common.config.ConfigResource;
+import org.apache.kafka.common.errors.ApiException;
import org.apache.kafka.common.errors.AuthenticationException;
import org.apache.kafka.common.errors.ClusterAuthorizationException;
import org.apache.kafka.common.errors.FencedInstanceIdException;
@@ -48,6 +50,7 @@
import org.apache.kafka.common.errors.InvalidRequestException;
import org.apache.kafka.common.errors.InvalidTopicException;
import org.apache.kafka.common.errors.LeaderNotAvailableException;
+import org.apache.kafka.common.errors.LogDirNotFoundException;
import org.apache.kafka.common.errors.TimeoutException;
import org.apache.kafka.common.errors.NotLeaderForPartitionException;
import org.apache.kafka.common.errors.OffsetOutOfRangeException;
@@ -59,6 +62,9 @@
import org.apache.kafka.common.errors.UnknownServerException;
import org.apache.kafka.common.errors.UnknownTopicOrPartitionException;
import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData;
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData;
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData.AlterReplicaLogDirPartitionResult;
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData.AlterReplicaLogDirTopicResult;
import org.apache.kafka.common.message.CreatePartitionsResponseData;
import org.apache.kafka.common.message.CreatePartitionsResponseData.CreatePartitionsTopicResult;
import org.apache.kafka.common.message.CreateAclsResponseData;
@@ -97,6 +103,7 @@
import org.apache.kafka.common.quota.ClientQuotaFilterComponent;
import org.apache.kafka.common.requests.AlterClientQuotasResponse;
import org.apache.kafka.common.requests.AlterPartitionReassignmentsResponse;
+import org.apache.kafka.common.requests.AlterReplicaLogDirsResponse;
import org.apache.kafka.common.requests.ApiError;
import org.apache.kafka.common.requests.CreateAclsResponse;
import org.apache.kafka.common.requests.CreatePartitionsResponse;
@@ -477,6 +484,21 @@ public void testCreateTopics() throws Exception {
}
}
+ @Test
+ public void testCreateTopicsPartialResponse() throws Exception {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ env.kafkaClient().setNodeApiVersions(NodeApiVersions.create());
+ env.kafkaClient().prepareResponse(body -> body instanceof CreateTopicsRequest,
+ prepareCreateTopicsResponse("myTopic", Errors.NONE));
+ CreateTopicsResult topicsResult = env.adminClient().createTopics(
+ asList(new NewTopic("myTopic", Collections.singletonMap(0, asList(0, 1, 2))),
+ new NewTopic("myTopic2", Collections.singletonMap(0, asList(0, 1, 2)))),
+ new CreateTopicsOptions().timeoutMs(10000));
+ topicsResult.values().get("myTopic").get();
+ TestUtils.assertFutureThrows(topicsResult.values().get("myTopic2"), ApiException.class);
+ }
+ }
+
@Test
public void testCreateTopicsRetryBackoff() throws Exception {
MockTime time = new MockTime();
@@ -570,6 +592,21 @@ public void testDeleteTopics() throws Exception {
}
}
+ @Test
+ public void testDeleteTopicsPartialResponse() throws Exception {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ env.kafkaClient().setNodeApiVersions(NodeApiVersions.create());
+
+ env.kafkaClient().prepareResponse(body -> body instanceof DeleteTopicsRequest,
+ prepareDeleteTopicsResponse("myTopic", Errors.NONE));
+ Map> values = env.adminClient().deleteTopics(asList("myTopic", "myOtherTopic"),
+ new DeleteTopicsOptions()).values();
+ values.get("myTopic").get();
+
+ TestUtils.assertFutureThrows(values.get("myOtherTopic"), ApiException.class);
+ }
+ }
+
@Test
public void testInvalidTopicNames() throws Exception {
try (AdminClientUnitTestEnv env = mockClientEnv()) {
@@ -3382,6 +3419,84 @@ public void testAlterClientQuotas() throws Exception {
}
}
+ @Test
+ public void testAlterReplicaLogDirsSuccess() throws Exception {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ createAlterLogDirsResponse(env, env.cluster().nodeById(0), Errors.NONE, 0);
+ createAlterLogDirsResponse(env, env.cluster().nodeById(1), Errors.NONE, 0);
+
+ TopicPartitionReplica tpr0 = new TopicPartitionReplica("topic", 0, 0);
+ TopicPartitionReplica tpr1 = new TopicPartitionReplica("topic", 0, 1);
+
+ Map logDirs = new HashMap<>();
+ logDirs.put(tpr0, "/data0");
+ logDirs.put(tpr1, "/data1");
+ AlterReplicaLogDirsResult result = env.adminClient().alterReplicaLogDirs(logDirs);
+ assertNull(result.values().get(tpr0).get());
+ assertNull(result.values().get(tpr1).get());
+ }
+ }
+
+ @Test
+ public void testAlterReplicaLogDirsLogDirNotFound() throws Exception {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ createAlterLogDirsResponse(env, env.cluster().nodeById(0), Errors.NONE, 0);
+ createAlterLogDirsResponse(env, env.cluster().nodeById(1), Errors.LOG_DIR_NOT_FOUND, 0);
+
+ TopicPartitionReplica tpr0 = new TopicPartitionReplica("topic", 0, 0);
+ TopicPartitionReplica tpr1 = new TopicPartitionReplica("topic", 0, 1);
+
+ Map logDirs = new HashMap<>();
+ logDirs.put(tpr0, "/data0");
+ logDirs.put(tpr1, "/data1");
+ AlterReplicaLogDirsResult result = env.adminClient().alterReplicaLogDirs(logDirs);
+ assertNull(result.values().get(tpr0).get());
+ TestUtils.assertFutureError(result.values().get(tpr1), LogDirNotFoundException.class);
+ }
+ }
+
+ @Test
+ public void testAlterReplicaLogDirsUnrequested() throws Exception {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ createAlterLogDirsResponse(env, env.cluster().nodeById(0), Errors.NONE, 1, 2);
+
+ TopicPartitionReplica tpr1 = new TopicPartitionReplica("topic", 1, 0);
+
+ Map logDirs = new HashMap<>();
+ logDirs.put(tpr1, "/data1");
+ AlterReplicaLogDirsResult result = env.adminClient().alterReplicaLogDirs(logDirs);
+ assertNull(result.values().get(tpr1).get());
+ }
+ }
+
+ @Test
+ public void testAlterReplicaLogDirsPartialResponse() throws Exception {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ createAlterLogDirsResponse(env, env.cluster().nodeById(0), Errors.NONE, 1);
+
+ TopicPartitionReplica tpr1 = new TopicPartitionReplica("topic", 1, 0);
+ TopicPartitionReplica tpr2 = new TopicPartitionReplica("topic", 2, 0);
+
+ Map logDirs = new HashMap<>();
+ logDirs.put(tpr1, "/data1");
+ logDirs.put(tpr2, "/data1");
+ AlterReplicaLogDirsResult result = env.adminClient().alterReplicaLogDirs(logDirs);
+ assertNull(result.values().get(tpr1).get());
+ TestUtils.assertFutureThrows(result.values().get(tpr2), ApiException.class);
+ }
+ }
+
+ private void createAlterLogDirsResponse(AdminClientUnitTestEnv env, Node node, Errors error, int... partitions) {
+ env.kafkaClient().prepareResponseFrom(new AlterReplicaLogDirsResponse(
+ new AlterReplicaLogDirsResponseData().setResults(singletonList(
+ new AlterReplicaLogDirTopicResult()
+ .setTopicName("topic")
+ .setPartitions(Arrays.stream(partitions).boxed().map(partitionId ->
+ new AlterReplicaLogDirPartitionResult()
+ .setPartitionIndex(partitionId)
+ .setErrorCode(error.code())).collect(Collectors.toList()))))), node);
+ }
+
private static MemberDescription convertToMemberDescriptions(DescribedGroupMember member,
MemberAssignment assignment) {
return new MemberDescription(member.memberId(),
diff --git a/clients/src/test/java/org/apache/kafka/common/requests/AlterReplicaLogDirsRequestTest.java b/clients/src/test/java/org/apache/kafka/common/requests/AlterReplicaLogDirsRequestTest.java
new file mode 100644
index 0000000000000..5ed919938a313
--- /dev/null
+++ b/clients/src/test/java/org/apache/kafka/common/requests/AlterReplicaLogDirsRequestTest.java
@@ -0,0 +1,88 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.common.requests;
+
+import java.util.HashMap;
+import java.util.Map;
+
+import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.errors.LogDirNotFoundException;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData.AlterReplicaLogDir;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData.AlterReplicaLogDirCollection;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData.AlterReplicaLogDirTopic;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData.AlterReplicaLogDirTopicCollection;
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData.AlterReplicaLogDirTopicResult;
+import org.apache.kafka.common.protocol.Errors;
+import org.junit.Test;
+
+import static java.util.Arrays.asList;
+import static java.util.Collections.singletonList;
+import static org.junit.Assert.assertEquals;
+
+public class AlterReplicaLogDirsRequestTest {
+
+ @Test
+ public void testErrorResponse() {
+ AlterReplicaLogDirsRequestData data = new AlterReplicaLogDirsRequestData()
+ .setDirs(new AlterReplicaLogDirCollection(
+ singletonList(new AlterReplicaLogDir()
+ .setPath("/data0")
+ .setTopics(new AlterReplicaLogDirTopicCollection(
+ singletonList(new AlterReplicaLogDirTopic()
+ .setName("topic")
+ .setPartitions(asList(0, 1, 2))).iterator()))).iterator()));
+ AlterReplicaLogDirsResponse errorResponse = new AlterReplicaLogDirsRequest.Builder(data).build()
+ .getErrorResponse(123, new LogDirNotFoundException("/data0"));
+ assertEquals(1, errorResponse.data().results().size());
+ AlterReplicaLogDirTopicResult topicResponse = errorResponse.data().results().get(0);
+ assertEquals("topic", topicResponse.topicName());
+ assertEquals(3, topicResponse.partitions().size());
+ for (int i = 0; i < 3; i++) {
+ assertEquals(i, topicResponse.partitions().get(i).partitionIndex());
+ assertEquals(Errors.LOG_DIR_NOT_FOUND.code(), topicResponse.partitions().get(i).errorCode());
+ }
+ }
+
+ @Test
+ public void testPartitionDir() {
+ AlterReplicaLogDirsRequestData data = new AlterReplicaLogDirsRequestData()
+ .setDirs(new AlterReplicaLogDirCollection(
+ asList(new AlterReplicaLogDir()
+ .setPath("/data0")
+ .setTopics(new AlterReplicaLogDirTopicCollection(
+ asList(new AlterReplicaLogDirTopic()
+ .setName("topic")
+ .setPartitions(asList(0, 1)),
+ new AlterReplicaLogDirTopic()
+ .setName("topic2")
+ .setPartitions(asList(7))).iterator())),
+ new AlterReplicaLogDir()
+ .setPath("/data1")
+ .setTopics(new AlterReplicaLogDirTopicCollection(
+ asList(new AlterReplicaLogDirTopic()
+ .setName("topic3")
+ .setPartitions(asList(12))).iterator()))).iterator()));
+ AlterReplicaLogDirsRequest request = new AlterReplicaLogDirsRequest.Builder(data).build();
+ Map expect = new HashMap<>();
+ expect.put(new TopicPartition("topic", 0), "/data0");
+ expect.put(new TopicPartition("topic", 1), "/data0");
+ expect.put(new TopicPartition("topic2", 7), "/data0");
+ expect.put(new TopicPartition("topic3", 12), "/data1");
+ assertEquals(expect, request.partitionDirs());
+ }
+}
diff --git a/clients/src/test/java/org/apache/kafka/common/requests/AlterReplicaLogDirsResponseTest.java b/clients/src/test/java/org/apache/kafka/common/requests/AlterReplicaLogDirsResponseTest.java
new file mode 100644
index 0000000000000..01826c69bc2ce
--- /dev/null
+++ b/clients/src/test/java/org/apache/kafka/common/requests/AlterReplicaLogDirsResponseTest.java
@@ -0,0 +1,57 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.common.requests;
+
+import java.util.Map;
+
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData;
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData.AlterReplicaLogDirPartitionResult;
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData.AlterReplicaLogDirTopicResult;
+import org.apache.kafka.common.protocol.Errors;
+import org.junit.Test;
+
+import static java.util.Arrays.asList;
+import static org.junit.Assert.assertEquals;
+
+public class AlterReplicaLogDirsResponseTest {
+
+ @Test
+ public void testErrorCounts() {
+ AlterReplicaLogDirsResponseData data = new AlterReplicaLogDirsResponseData()
+ .setResults(asList(
+ new AlterReplicaLogDirTopicResult()
+ .setTopicName("t0")
+ .setPartitions(asList(
+ new AlterReplicaLogDirPartitionResult()
+ .setPartitionIndex(0)
+ .setErrorCode(Errors.LOG_DIR_NOT_FOUND.code()),
+ new AlterReplicaLogDirPartitionResult()
+ .setPartitionIndex(1)
+ .setErrorCode(Errors.NONE.code()))),
+ new AlterReplicaLogDirTopicResult()
+ .setTopicName("t1")
+ .setPartitions(asList(
+ new AlterReplicaLogDirPartitionResult()
+ .setPartitionIndex(0)
+ .setErrorCode(Errors.LOG_DIR_NOT_FOUND.code())))));
+ Map counts = new AlterReplicaLogDirsResponse(data).errorCounts();
+ assertEquals(2, counts.size());
+ assertEquals(Integer.valueOf(2), counts.get(Errors.LOG_DIR_NOT_FOUND));
+ assertEquals(Integer.valueOf(1), counts.get(Errors.NONE));
+
+ }
+}
diff --git a/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java b/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java
index 69b658a16bd99..efae74e326671 100644
--- a/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java
+++ b/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java
@@ -38,6 +38,10 @@
import org.apache.kafka.common.message.AlterConfigsResponseData;
import org.apache.kafka.common.message.AlterPartitionReassignmentsRequestData;
import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData.AlterReplicaLogDirTopic;
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData.AlterReplicaLogDirTopicCollection;
+import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData;
import org.apache.kafka.common.message.ApiVersionsRequestData;
import org.apache.kafka.common.message.ApiVersionsResponseData;
import org.apache.kafka.common.message.ApiVersionsResponseData.ApiVersionsResponseKey;
@@ -464,6 +468,9 @@ public void testSerialization() throws Exception {
checkRequest(createOffsetDeleteRequest(), true);
checkErrorResponse(createOffsetDeleteRequest(), new UnknownServerException(), true);
checkResponse(createOffsetDeleteResponse(), 0, true);
+ checkRequest(createAlterReplicaLogDirsRequest(), true);
+ checkErrorResponse(createAlterReplicaLogDirsRequest(), new UnknownServerException(), true);
+ checkResponse(createAlterReplicaLogDirsResponse(), 0, true);
}
@Test
@@ -2198,4 +2205,34 @@ private OffsetDeleteResponse createOffsetDeleteResponse() {
return new OffsetDeleteResponse(data);
}
+ private AlterReplicaLogDirsRequest createAlterReplicaLogDirsRequest() {
+ AlterReplicaLogDirsRequestData data = new AlterReplicaLogDirsRequestData();
+ data.dirs().add(
+ new AlterReplicaLogDirsRequestData.AlterReplicaLogDir()
+ .setPath("/data0")
+ .setTopics(new AlterReplicaLogDirTopicCollection(Collections.singletonList(
+ new AlterReplicaLogDirTopic()
+ .setPartitions(singletonList(0))
+ .setName("topic")
+ ).iterator())
+ )
+ );
+ return new AlterReplicaLogDirsRequest.Builder(data).build((short) 0);
+ }
+
+ private AlterReplicaLogDirsResponse createAlterReplicaLogDirsResponse() {
+ AlterReplicaLogDirsResponseData data = new AlterReplicaLogDirsResponseData();
+ data.results().add(
+ new AlterReplicaLogDirsResponseData.AlterReplicaLogDirTopicResult()
+ .setTopicName("topic")
+ .setPartitions(Collections.singletonList(
+ new AlterReplicaLogDirsResponseData.AlterReplicaLogDirPartitionResult()
+ .setPartitionIndex(0)
+ .setErrorCode(Errors.LOG_DIR_NOT_FOUND.code())
+ )
+ )
+ );
+ return new AlterReplicaLogDirsResponse(data);
+ }
+
}
diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala
index cf1cb98925944..152e6039553d2 100644
--- a/core/src/main/scala/kafka/server/KafkaApis.scala
+++ b/core/src/main/scala/kafka/server/KafkaApis.scala
@@ -51,7 +51,7 @@ import org.apache.kafka.common.internals.Topic.{GROUP_METADATA_TOPIC_NAME, TRANS
import org.apache.kafka.common.message.AlterConfigsResponseData.AlterConfigsResourceResponse
import org.apache.kafka.common.message.CreateTopicsRequestData.CreatableTopic
import org.apache.kafka.common.message.CreatePartitionsResponseData.CreatePartitionsTopicResult
-import org.apache.kafka.common.message.{AddOffsetsToTxnResponseData, AlterConfigsResponseData, AlterPartitionReassignmentsResponseData, CreateAclsResponseData, CreatePartitionsResponseData, CreateTopicsResponseData, DeleteAclsResponseData, DeleteGroupsResponseData, DeleteRecordsResponseData, DeleteTopicsResponseData, DescribeAclsResponseData, DescribeGroupsResponseData, DescribeLogDirsResponseData, EndTxnResponseData, ExpireDelegationTokenResponseData, FindCoordinatorResponseData, HeartbeatResponseData, InitProducerIdResponseData, JoinGroupResponseData, LeaveGroupResponseData, ListGroupsResponseData, ListPartitionReassignmentsResponseData, OffsetCommitRequestData, OffsetCommitResponseData, OffsetDeleteResponseData, RenewDelegationTokenResponseData, SaslAuthenticateResponseData, SaslHandshakeResponseData, StopReplicaResponseData, SyncGroupResponseData, UpdateMetadataResponseData}
+import org.apache.kafka.common.message.{AddOffsetsToTxnResponseData, AlterConfigsResponseData, AlterPartitionReassignmentsResponseData, AlterReplicaLogDirsResponseData, CreateAclsResponseData, CreatePartitionsResponseData, CreateTopicsResponseData, DeleteAclsResponseData, DeleteGroupsResponseData, DeleteRecordsResponseData, DeleteTopicsResponseData, DescribeAclsResponseData, DescribeGroupsResponseData, DescribeLogDirsResponseData, EndTxnResponseData, ExpireDelegationTokenResponseData, FindCoordinatorResponseData, HeartbeatResponseData, InitProducerIdResponseData, JoinGroupResponseData, LeaveGroupResponseData, ListGroupsResponseData, ListPartitionReassignmentsResponseData, OffsetCommitRequestData, OffsetCommitResponseData, OffsetDeleteResponseData, RenewDelegationTokenResponseData, SaslAuthenticateResponseData, SaslHandshakeResponseData, StopReplicaResponseData, SyncGroupResponseData, UpdateMetadataResponseData}
import org.apache.kafka.common.message.CreateTopicsResponseData.{CreatableTopicResult, CreatableTopicResultCollection}
import org.apache.kafka.common.message.DeleteGroupsResponseData.{DeletableGroupResult, DeletableGroupResultCollection}
import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.{ReassignablePartitionResponse, ReassignableTopicResponse}
@@ -2590,13 +2590,24 @@ class KafkaApis(val requestChannel: RequestChannel,
def handleAlterReplicaLogDirsRequest(request: RequestChannel.Request): Unit = {
val alterReplicaDirsRequest = request.body[AlterReplicaLogDirsRequest]
- val responseMap = {
- if (authorize(request.context, ALTER, CLUSTER, CLUSTER_NAME))
- replicaManager.alterReplicaLogDirs(alterReplicaDirsRequest.partitionDirs.asScala)
- else
- alterReplicaDirsRequest.partitionDirs.asScala.keys.map((_, Errors.CLUSTER_AUTHORIZATION_FAILED)).toMap
+ if (authorize(request.context, ALTER, CLUSTER, CLUSTER_NAME)) {
+ val result = replicaManager.alterReplicaLogDirs(alterReplicaDirsRequest.partitionDirs.asScala)
+ sendResponseMaybeThrottle(request, requestThrottleMs =>
+ new AlterReplicaLogDirsResponse(new AlterReplicaLogDirsResponseData()
+ .setResults(result.groupBy(_._1.topic).map {
+ case (topic, errors) => new AlterReplicaLogDirsResponseData.AlterReplicaLogDirTopicResult()
+ .setTopicName(topic)
+ .setPartitions(errors.map {
+ case (tp, error) => new AlterReplicaLogDirsResponseData.AlterReplicaLogDirPartitionResult()
+ .setPartitionIndex(tp.partition)
+ .setErrorCode(error.code)
+ }.toList.asJava)
+ }.toList.asJava)
+ .setThrottleTimeMs(requestThrottleMs)))
+ } else {
+ sendResponseMaybeThrottle(request, requestThrottleMs =>
+ alterReplicaDirsRequest.getErrorResponse(requestThrottleMs, Errors.CLUSTER_AUTHORIZATION_FAILED.exception))
}
- sendResponseMaybeThrottle(request, requestThrottleMs => new AlterReplicaLogDirsResponse(requestThrottleMs, responseMap.asJava))
}
def handleDescribeLogDirsRequest(request: RequestChannel.Request): Unit = {
diff --git a/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala
index 904a2be44c505..f3b259b634a3f 100644
--- a/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala
+++ b/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala
@@ -44,7 +44,7 @@ import org.apache.kafka.common.message.LeaderAndIsrRequestData.LeaderAndIsrParti
import org.apache.kafka.common.message.LeaveGroupRequestData.MemberIdentity
import org.apache.kafka.common.message.StopReplicaRequestData.{StopReplicaPartitionState, StopReplicaTopicState}
import org.apache.kafka.common.message.UpdateMetadataRequestData.{UpdateMetadataBroker, UpdateMetadataEndpoint, UpdateMetadataPartitionState}
-import org.apache.kafka.common.message.{AlterPartitionReassignmentsRequestData, ControlledShutdownRequestData, CreateAclsRequestData, CreatePartitionsRequestData, CreateTopicsRequestData, DeleteAclsRequestData, DeleteGroupsRequestData, DeleteRecordsRequestData, DeleteTopicsRequestData, DescribeGroupsRequestData, DescribeLogDirsRequestData, FindCoordinatorRequestData, HeartbeatRequestData, IncrementalAlterConfigsRequestData, JoinGroupRequestData, ListPartitionReassignmentsRequestData, OffsetCommitRequestData, SyncGroupRequestData}
+import org.apache.kafka.common.message.{AlterPartitionReassignmentsRequestData, AlterReplicaLogDirsRequestData, ControlledShutdownRequestData, CreateAclsRequestData, CreatePartitionsRequestData, CreateTopicsRequestData, DeleteAclsRequestData, DeleteGroupsRequestData, DeleteRecordsRequestData, DeleteTopicsRequestData, DescribeGroupsRequestData, DescribeLogDirsRequestData, FindCoordinatorRequestData, HeartbeatRequestData, IncrementalAlterConfigsRequestData, JoinGroupRequestData, ListPartitionReassignmentsRequestData, OffsetCommitRequestData, SyncGroupRequestData}
import org.apache.kafka.common.network.ListenerName
import org.apache.kafka.common.protocol.{ApiKeys, Errors}
import org.apache.kafka.common.record.{CompressionType, MemoryRecords, RecordBatch, Records, SimpleRecord}
@@ -193,7 +193,9 @@ class AuthorizerIntegrationTest extends BaseRequestTest {
ApiKeys.CREATE_ACLS -> ((resp: CreateAclsResponse) => Errors.forCode(resp.results.asScala.head.errorCode)),
ApiKeys.DESCRIBE_ACLS -> ((resp: DescribeAclsResponse) => resp.error.error),
ApiKeys.DELETE_ACLS -> ((resp: DeleteAclsResponse) => Errors.forCode(resp.filterResults.asScala.head.errorCode)),
- ApiKeys.ALTER_REPLICA_LOG_DIRS -> ((resp: AlterReplicaLogDirsResponse) => resp.responses.get(tp)),
+ ApiKeys.ALTER_REPLICA_LOG_DIRS -> ((resp: AlterReplicaLogDirsResponse) => Errors.forCode(resp.data.results.asScala
+ .find(x => x.topicName == tp.topic).get.partitions.asScala
+ .find(p => p.partitionIndex == tp.partition).get.errorCode)),
ApiKeys.DESCRIBE_LOG_DIRS -> ((resp: DescribeLogDirsResponse) =>
if (resp.logDirInfos.size() > 0) resp.logDirInfos.asScala.head._2.error else Errors.CLUSTER_AUTHORIZATION_FAILED),
ApiKeys.CREATE_PARTITIONS -> ((resp: CreatePartitionsResponse) => Errors.forCode(resp.data.results.asScala.head.errorCode())),
@@ -537,7 +539,16 @@ class AuthorizerIntegrationTest extends BaseRequestTest {
.setPermissionType(AclPermissionType.DENY.code)))
).build()
- private def alterReplicaLogDirsRequest = new AlterReplicaLogDirsRequest.Builder(Collections.singletonMap(tp, logDir)).build()
+ private def alterReplicaLogDirsRequest = {
+ val dir = new AlterReplicaLogDirsRequestData.AlterReplicaLogDir()
+ .setPath(logDir)
+ dir.topics.add(new AlterReplicaLogDirsRequestData.AlterReplicaLogDirTopic()
+ .setName(tp.topic)
+ .setPartitions(Collections.singletonList(tp.partition)))
+ val data = new AlterReplicaLogDirsRequestData();
+ data.dirs.add(dir)
+ new AlterReplicaLogDirsRequest.Builder(data).build()
+ }
private def describeLogDirsRequest = new DescribeLogDirsRequest.Builder(new DescribeLogDirsRequestData().setTopics(new DescribeLogDirsRequestData.DescribableLogDirTopicCollection(Collections.singleton(
new DescribeLogDirsRequestData.DescribableLogDirTopic().setTopic(tp.topic).setPartitionIndex(Collections.singletonList(tp.partition))).iterator()))).build()
diff --git a/core/src/test/scala/unit/kafka/server/AlterReplicaLogDirsRequestTest.scala b/core/src/test/scala/unit/kafka/server/AlterReplicaLogDirsRequestTest.scala
index 236e60ed6ba83..425456f9b6cb4 100644
--- a/core/src/test/scala/unit/kafka/server/AlterReplicaLogDirsRequestTest.scala
+++ b/core/src/test/scala/unit/kafka/server/AlterReplicaLogDirsRequestTest.scala
@@ -21,6 +21,7 @@ import java.io.File
import kafka.utils._
import org.apache.kafka.common.TopicPartition
+import org.apache.kafka.common.message.AlterReplicaLogDirsRequestData
import org.apache.kafka.common.protocol.Errors
import org.apache.kafka.common.requests.{AlterReplicaLogDirsRequest, AlterReplicaLogDirsResponse}
import org.junit.Assert._
@@ -36,6 +37,12 @@ class AlterReplicaLogDirsRequestTest extends BaseRequestTest {
val topic = "topic"
+ private def findErrorForPartition(response: AlterReplicaLogDirsResponse, tp: TopicPartition): Errors = {
+ Errors.forCode(response.data.results.asScala
+ .find(x => x.topicName == tp.topic).get.partitions.asScala
+ .find(p => p.partitionIndex == tp.partition).get.errorCode)
+ }
+
@Test
def testAlterReplicaLogDirsRequest(): Unit = {
val partitionNum = 5
@@ -48,7 +55,7 @@ class AlterReplicaLogDirsRequestTest extends BaseRequestTest {
// The response should show error UNKNOWN_TOPIC_OR_PARTITION for all partitions
(0 until partitionNum).foreach { partition =>
val tp = new TopicPartition(topic, partition)
- assertEquals(Errors.UNKNOWN_TOPIC_OR_PARTITION, alterReplicaLogDirsResponse1.responses().get(tp))
+ assertEquals(Errors.UNKNOWN_TOPIC_OR_PARTITION, findErrorForPartition(alterReplicaLogDirsResponse1, tp))
assertTrue(servers.head.logManager.getLog(tp).isEmpty)
}
@@ -64,7 +71,7 @@ class AlterReplicaLogDirsRequestTest extends BaseRequestTest {
// The response should succeed for all partitions
(0 until partitionNum).foreach { partition =>
val tp = new TopicPartition(topic, partition)
- assertEquals(Errors.NONE, alterReplicaLogDirsResponse2.responses().get(tp))
+ assertEquals(Errors.NONE, findErrorForPartition(alterReplicaLogDirsResponse2, tp))
TestUtils.waitUntilTrue(() => {
logDir2 == servers.head.logManager.getLog(new TopicPartition(topic, partition)).get.dir.getParent
}, "timed out waiting for replica movement")
@@ -83,8 +90,8 @@ class AlterReplicaLogDirsRequestTest extends BaseRequestTest {
partitionDirs1.put(new TopicPartition(topic, 0), "invalidDir")
partitionDirs1.put(new TopicPartition(topic, 1), validDir1)
val alterReplicaDirResponse1 = sendAlterReplicaLogDirsRequest(partitionDirs1.toMap)
- assertEquals(Errors.LOG_DIR_NOT_FOUND, alterReplicaDirResponse1.responses().get(new TopicPartition(topic, 0)))
- assertEquals(Errors.UNKNOWN_TOPIC_OR_PARTITION, alterReplicaDirResponse1.responses().get(new TopicPartition(topic, 1)))
+ assertEquals(Errors.LOG_DIR_NOT_FOUND, findErrorForPartition(alterReplicaDirResponse1, new TopicPartition(topic, 0)))
+ assertEquals(Errors.UNKNOWN_TOPIC_OR_PARTITION, findErrorForPartition(alterReplicaDirResponse1, new TopicPartition(topic, 1)))
createTopic(topic, 3, 1)
@@ -93,8 +100,8 @@ class AlterReplicaLogDirsRequestTest extends BaseRequestTest {
partitionDirs2.put(new TopicPartition(topic, 0), "invalidDir")
partitionDirs2.put(new TopicPartition(topic, 1), validDir2)
val alterReplicaDirResponse2 = sendAlterReplicaLogDirsRequest(partitionDirs2.toMap)
- assertEquals(Errors.LOG_DIR_NOT_FOUND, alterReplicaDirResponse2.responses().get(new TopicPartition(topic, 0)))
- assertEquals(Errors.NONE, alterReplicaDirResponse2.responses().get(new TopicPartition(topic, 1)))
+ assertEquals(Errors.LOG_DIR_NOT_FOUND, findErrorForPartition(alterReplicaDirResponse2, new TopicPartition(topic, 0)))
+ assertEquals(Errors.NONE, findErrorForPartition(alterReplicaDirResponse2, new TopicPartition(topic, 1)))
// Test AlterReplicaDirRequest after topic creation and log directory failure
servers.head.logDirFailureChannel.maybeAddOfflineLogDir(offlineDir, "", new java.io.IOException())
@@ -104,13 +111,26 @@ class AlterReplicaLogDirsRequestTest extends BaseRequestTest {
partitionDirs3.put(new TopicPartition(topic, 1), validDir3)
partitionDirs3.put(new TopicPartition(topic, 2), offlineDir)
val alterReplicaDirResponse3 = sendAlterReplicaLogDirsRequest(partitionDirs3.toMap)
- assertEquals(Errors.LOG_DIR_NOT_FOUND, alterReplicaDirResponse3.responses().get(new TopicPartition(topic, 0)))
- assertEquals(Errors.KAFKA_STORAGE_ERROR, alterReplicaDirResponse3.responses().get(new TopicPartition(topic, 1)))
- assertEquals(Errors.KAFKA_STORAGE_ERROR, alterReplicaDirResponse3.responses().get(new TopicPartition(topic, 2)))
+ assertEquals(Errors.LOG_DIR_NOT_FOUND, findErrorForPartition(alterReplicaDirResponse3, new TopicPartition(topic, 0)))
+ assertEquals(Errors.KAFKA_STORAGE_ERROR, findErrorForPartition(alterReplicaDirResponse3, new TopicPartition(topic, 1)))
+ assertEquals(Errors.KAFKA_STORAGE_ERROR, findErrorForPartition(alterReplicaDirResponse3, new TopicPartition(topic, 2)))
}
private def sendAlterReplicaLogDirsRequest(partitionDirs: Map[TopicPartition, String]): AlterReplicaLogDirsResponse = {
- val request = new AlterReplicaLogDirsRequest.Builder(partitionDirs.asJava).build()
+ val logDirs = partitionDirs.groupBy{case (_, dir) => dir}.map{ case(dir, tps) =>
+ new AlterReplicaLogDirsRequestData.AlterReplicaLogDir()
+ .setPath(dir)
+ .setTopics(new AlterReplicaLogDirsRequestData.AlterReplicaLogDirTopicCollection(
+ tps.groupBy { case (tp, _) => tp.topic }
+ .map { case (topic, tpPartitions) =>
+ new AlterReplicaLogDirsRequestData.AlterReplicaLogDirTopic()
+ .setName(topic)
+ .setPartitions(tpPartitions.map{case (tp, _) => tp.partition.asInstanceOf[Integer]}.toList.asJava)
+ }.toList.asJava.iterator))
+ }
+ val data = new AlterReplicaLogDirsRequestData()
+ .setDirs(new AlterReplicaLogDirsRequestData.AlterReplicaLogDirCollection(logDirs.asJava.iterator))
+ val request = new AlterReplicaLogDirsRequest.Builder(data).build()
connectAndReceive[AlterReplicaLogDirsResponse](request, destination = controllerSocketServer)
}
diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
index 1b6894b5a8783..9cf322b0505a0 100644
--- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
+++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
@@ -1707,4 +1707,47 @@ class KafkaApisTest {
0, 0, partitionStates.asJava, Seq(broker).asJava).build()
metadataCache.updateMetadata(correlationId = 0, updateMetadataRequest)
}
+
+ @Test
+ def testAlterReplicaLogDirs(): Unit = {
+ val data = new AlterReplicaLogDirsRequestData()
+ val dir = new AlterReplicaLogDirsRequestData.AlterReplicaLogDir()
+ .setPath("/foo")
+ dir.topics().add(new AlterReplicaLogDirsRequestData.AlterReplicaLogDirTopic().setName("t0").setPartitions(asList(0, 1, 2)))
+ data.dirs().add(dir)
+ val alterReplicaLogDirsRequest = new AlterReplicaLogDirsRequest.Builder(
+ data
+ ).build()
+ val request = buildRequest(alterReplicaLogDirsRequest)
+
+ EasyMock.reset(replicaManager, clientRequestQuotaManager, requestChannel)
+
+ val capturedResponse = expectNoThrottling()
+ val t0p0 = new TopicPartition("t0", 0)
+ val t0p1 = new TopicPartition("t0", 1)
+ val t0p2 = new TopicPartition("t0", 2)
+ val partitionResults = Map(
+ t0p0 -> Errors.NONE,
+ t0p1 -> Errors.LOG_DIR_NOT_FOUND,
+ t0p2 -> Errors.INVALID_TOPIC_EXCEPTION)
+ EasyMock.expect(replicaManager.alterReplicaLogDirs(EasyMock.eq(Map(
+ t0p0 -> "/foo",
+ t0p1 -> "/foo",
+ t0p2 -> "/foo"))))
+ .andReturn(partitionResults)
+ EasyMock.replay(replicaManager, clientQuotaManager, clientRequestQuotaManager, requestChannel)
+
+ createKafkaApis().handleAlterReplicaLogDirsRequest(request)
+
+ val response = readResponse(ApiKeys.ALTER_REPLICA_LOG_DIRS, alterReplicaLogDirsRequest, capturedResponse)
+ .asInstanceOf[AlterReplicaLogDirsResponse]
+ assertEquals(partitionResults, response.data.results.asScala.flatMap { tr =>
+ tr.partitions().asScala.map { pr =>
+ new TopicPartition(tr.topicName, pr.partitionIndex) -> Errors.forCode(pr.errorCode)
+ }
+ }.toMap)
+ assertEquals(Map(Errors.NONE -> 1,
+ Errors.LOG_DIR_NOT_FOUND -> 1,
+ Errors.INVALID_TOPIC_EXCEPTION -> 1).asJava, response.errorCounts)
+ }
}
diff --git a/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala b/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala
index 6dd151828dba8..0e4ef315e02ba 100644
--- a/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala
+++ b/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala
@@ -468,7 +468,14 @@ class RequestQuotaTest extends BaseRequestTest {
))), true)
case ApiKeys.ALTER_REPLICA_LOG_DIRS =>
- new AlterReplicaLogDirsRequest.Builder(Collections.singletonMap(tp, logDir))
+ val dir = new AlterReplicaLogDirsRequestData.AlterReplicaLogDir()
+ .setPath(logDir)
+ dir.topics.add(new AlterReplicaLogDirsRequestData.AlterReplicaLogDirTopic()
+ .setName(tp.topic)
+ .setPartitions(Collections.singletonList(tp.partition)))
+ val data = new AlterReplicaLogDirsRequestData();
+ data.dirs.add(dir)
+ new AlterReplicaLogDirsRequest.Builder(data)
case ApiKeys.DESCRIBE_LOG_DIRS =>
val data = new DescribeLogDirsRequestData()