diff --git a/checkstyle/import-control.xml b/checkstyle/import-control.xml
index ceee02adffa50..5368ef217bfbf 100644
--- a/checkstyle/import-control.xml
+++ b/checkstyle/import-control.xml
@@ -318,6 +318,7 @@
+
@@ -329,6 +330,12 @@
+
+
+
+
+
+
diff --git a/clients/src/main/java/org/apache/kafka/clients/ClientRequest.java b/clients/src/main/java/org/apache/kafka/clients/ClientRequest.java
index cf3b911495cb5..abba79567eb1e 100644
--- a/clients/src/main/java/org/apache/kafka/clients/ClientRequest.java
+++ b/clients/src/main/java/org/apache/kafka/clients/ClientRequest.java
@@ -83,14 +83,14 @@ public ApiKeys apiKey() {
}
public RequestHeader makeHeader(short version) {
- short requestApiKey = requestBuilder.apiKey().id;
+ ApiKeys requestApiKey = apiKey();
return new RequestHeader(
new RequestHeaderData()
- .setRequestApiKey(requestApiKey)
+ .setRequestApiKey(requestApiKey.id)
.setRequestApiVersion(version)
.setClientId(clientId)
.setCorrelationId(correlationId),
- ApiKeys.forId(requestApiKey).requestHeaderVersion(version));
+ requestApiKey.requestHeaderVersion(version));
}
public AbstractRequest.Builder> requestBuilder() {
diff --git a/clients/src/main/java/org/apache/kafka/clients/NetworkClient.java b/clients/src/main/java/org/apache/kafka/clients/NetworkClient.java
index 7ec3b93be4b2e..99bfd85d98530 100644
--- a/clients/src/main/java/org/apache/kafka/clients/NetworkClient.java
+++ b/clients/src/main/java/org/apache/kafka/clients/NetworkClient.java
@@ -870,9 +870,8 @@ private void handleCompletedReceives(List responses, long now) {
InFlightRequest req = inFlightRequests.completeNext(source);
AbstractResponse response = parseResponse(receive.payload(), req.header);
- if (throttleTimeSensor != null) {
- throttleTimeSensor.record(response.throttleTimeMs());
- }
+ if (throttleTimeSensor != null)
+ throttleTimeSensor.record(response.throttleTimeMs(), now);
if (log.isDebugEnabled()) {
log.debug("Received {} response from node {} for request with header {}: {}",
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 9c931cd7556f7..cf0a6f6949f90 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
@@ -1182,8 +1182,7 @@ private void handleResponses(long now, List responses) {
try {
call.handleResponse(response.responseBody());
if (log.isTraceEnabled())
- log.trace("{} got response {}", call,
- response.responseBody().toString(response.requestHeader().apiVersion()));
+ log.trace("{} got response {}", call, response.responseBody());
} catch (Throwable t) {
if (log.isTraceEnabled())
log.trace("{} handleResponse failed with {}", call, prettyPrintException(t));
diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinator.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinator.java
index af6e262078d7a..4ca71bb626a06 100644
--- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinator.java
+++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinator.java
@@ -465,7 +465,13 @@ boolean joinGroupIfNeeded(final Timer timer) {
}
} else {
final RuntimeException exception = future.exception();
- log.info("Rebalance failed.", exception);
+
+ // we do not need to log error for memberId required,
+ // since it is not really an error and is transient
+ if (!(exception instanceof MemberIdRequiredException)) {
+ log.info("Rebalance failed.", exception);
+ }
+
resetJoinGroupFuture();
if (exception instanceof UnknownMemberIdException ||
exception instanceof RebalanceInProgressException ||
diff --git a/clients/src/main/java/org/apache/kafka/common/protocol/SendBuilder.java b/clients/src/main/java/org/apache/kafka/common/protocol/SendBuilder.java
index 19d5600da89b8..6ed4c19a5719b 100644
--- a/clients/src/main/java/org/apache/kafka/common/protocol/SendBuilder.java
+++ b/clients/src/main/java/org/apache/kafka/common/protocol/SendBuilder.java
@@ -22,6 +22,7 @@
import org.apache.kafka.common.record.MemoryRecords;
import org.apache.kafka.common.record.MultiRecordsSend;
import org.apache.kafka.common.requests.RequestHeader;
+import org.apache.kafka.common.requests.RequestUtils;
import org.apache.kafka.common.requests.ResponseHeader;
import org.apache.kafka.common.utils.ByteUtils;
@@ -181,7 +182,7 @@ public Send build() {
public static Send buildRequestSend(
String destination,
RequestHeader header,
- ApiMessage apiRequest
+ Message apiRequest
) {
return buildSend(
destination,
@@ -195,7 +196,7 @@ public static Send buildRequestSend(
public static Send buildResponseSend(
String destination,
ResponseHeader header,
- ApiMessage apiResponse,
+ Message apiResponse,
short apiVersion
) {
return buildSend(
@@ -209,16 +210,13 @@ public static Send buildResponseSend(
private static Send buildSend(
String destination,
- ApiMessage header,
+ Message header,
short headerVersion,
- ApiMessage apiMessage,
+ Message apiMessage,
short apiVersion
) {
ObjectSerializationCache serializationCache = new ObjectSerializationCache();
- MessageSizeAccumulator messageSize = new MessageSizeAccumulator();
-
- header.addSize(messageSize, serializationCache, headerVersion);
- apiMessage.addSize(messageSize, serializationCache, apiVersion);
+ MessageSizeAccumulator messageSize = RequestUtils.size(serializationCache, header, headerVersion, apiMessage, apiVersion);
int totalSize = messageSize.totalSize();
int sizeExcludingZeroCopyFields = totalSize - messageSize.zeroCopySize();
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequest.java
index a0121666dfcf2..e86cd54224ad7 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequest.java
@@ -17,14 +17,12 @@
package org.apache.kafka.common.requests;
import org.apache.kafka.common.errors.UnsupportedVersionException;
-import org.apache.kafka.common.message.FetchRequestData;
-import org.apache.kafka.common.message.AlterIsrRequestData;
-import org.apache.kafka.common.message.ProduceRequestData;
-import org.apache.kafka.common.network.NetworkSend;
import org.apache.kafka.common.network.Send;
import org.apache.kafka.common.protocol.ApiKeys;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.Message;
+import org.apache.kafka.common.protocol.ObjectSerializationCache;
+import org.apache.kafka.common.protocol.SendBuilder;
import java.nio.ByteBuffer;
import java.util.Map;
@@ -79,13 +77,13 @@ public T build() {
}
private final short version;
- public final ApiKeys api;
+ private final ApiKeys apiKey;
- public AbstractRequest(ApiKeys api, short version) {
- if (!api.isVersionSupported(version))
- throw new UnsupportedVersionException("The " + api + " protocol does not support version " + version);
+ public AbstractRequest(ApiKeys apiKey, short version) {
+ if (!apiKey.isVersionSupported(version))
+ throw new UnsupportedVersionException("The " + apiKey + " protocol does not support version " + version);
this.version = version;
- this.api = api;
+ this.apiKey = apiKey;
}
/**
@@ -95,21 +93,33 @@ public short version() {
return version;
}
- public Send toSend(String destination, RequestHeader header) {
- return new NetworkSend(destination, serialize(header));
+ public ApiKeys apiKey() {
+ return apiKey;
}
- /**
- * Use with care, typically {@link #toSend(String, RequestHeader)} should be used instead.
- */
- public ByteBuffer serialize(RequestHeader header) {
- return RequestUtils.serialize(header.toStruct(), toStruct());
+ public final Send toSend(String destination, RequestHeader header) {
+ return SendBuilder.buildRequestSend(destination, header, data());
+ }
+
+ // Visible for testing
+ public final ByteBuffer serializeWithHeader(RequestHeader header) {
+ return RequestUtils.serialize(header.data(), header.headerVersion(), data(), version);
+ }
+
+ protected abstract Message data();
+
+ // Visible for testing
+ public final ByteBuffer serializeBody() {
+ return RequestUtils.serialize(null, (short) 0, data(), version);
}
- protected abstract Struct toStruct();
+ // Visible for testing
+ final int sizeInBytes() {
+ return data().size(new ObjectSerializationCache(), version);
+ }
public String toString(boolean verbose) {
- return toStruct().toString();
+ return data().toString();
}
@Override
@@ -144,126 +154,131 @@ public Map errorCounts(Throwable e) {
/**
* Factory method for getting a request object based on ApiKey ID and a version
*/
- public static AbstractRequest parseRequest(ApiKeys apiKey, short apiVersion, Struct struct) {
+ public static RequestAndSize parseRequest(ApiKeys apiKey, short apiVersion, ByteBuffer buffer) {
+ int bufferSize = buffer.remaining();
+ return new RequestAndSize(doParseRequest(apiKey, apiVersion, buffer), bufferSize);
+ }
+
+ private static AbstractRequest doParseRequest(ApiKeys apiKey, short apiVersion, ByteBuffer buffer) {
switch (apiKey) {
case PRODUCE:
- return new ProduceRequest(new ProduceRequestData(struct, apiVersion), apiVersion);
+ return ProduceRequest.parse(buffer, apiVersion);
case FETCH:
- return new FetchRequest(new FetchRequestData(struct, apiVersion), apiVersion);
+ return FetchRequest.parse(buffer, apiVersion);
case LIST_OFFSETS:
- return new ListOffsetRequest(struct, apiVersion);
+ return ListOffsetRequest.parse(buffer, apiVersion);
case METADATA:
- return new MetadataRequest(struct, apiVersion);
+ return MetadataRequest.parse(buffer, apiVersion);
case OFFSET_COMMIT:
- return new OffsetCommitRequest(struct, apiVersion);
+ return OffsetCommitRequest.parse(buffer, apiVersion);
case OFFSET_FETCH:
- return new OffsetFetchRequest(struct, apiVersion);
+ return OffsetFetchRequest.parse(buffer, apiVersion);
case FIND_COORDINATOR:
- return new FindCoordinatorRequest(struct, apiVersion);
+ return FindCoordinatorRequest.parse(buffer, apiVersion);
case JOIN_GROUP:
- return new JoinGroupRequest(struct, apiVersion);
+ return JoinGroupRequest.parse(buffer, apiVersion);
case HEARTBEAT:
- return new HeartbeatRequest(struct, apiVersion);
+ return HeartbeatRequest.parse(buffer, apiVersion);
case LEAVE_GROUP:
- return new LeaveGroupRequest(struct, apiVersion);
+ return LeaveGroupRequest.parse(buffer, apiVersion);
case SYNC_GROUP:
- return new SyncGroupRequest(struct, apiVersion);
+ return SyncGroupRequest.parse(buffer, apiVersion);
case STOP_REPLICA:
- return new StopReplicaRequest(struct, apiVersion);
+ return StopReplicaRequest.parse(buffer, apiVersion);
case CONTROLLED_SHUTDOWN:
- return new ControlledShutdownRequest(struct, apiVersion);
+ return ControlledShutdownRequest.parse(buffer, apiVersion);
case UPDATE_METADATA:
- return new UpdateMetadataRequest(struct, apiVersion);
+ return UpdateMetadataRequest.parse(buffer, apiVersion);
case LEADER_AND_ISR:
- return new LeaderAndIsrRequest(struct, apiVersion);
+ return LeaderAndIsrRequest.parse(buffer, apiVersion);
case DESCRIBE_GROUPS:
- return new DescribeGroupsRequest(struct, apiVersion);
+ return DescribeGroupsRequest.parse(buffer, apiVersion);
case LIST_GROUPS:
- return new ListGroupsRequest(struct, apiVersion);
+ return ListGroupsRequest.parse(buffer, apiVersion);
case SASL_HANDSHAKE:
- return new SaslHandshakeRequest(struct, apiVersion);
+ return SaslHandshakeRequest.parse(buffer, apiVersion);
case API_VERSIONS:
- return new ApiVersionsRequest(struct, apiVersion);
+ return ApiVersionsRequest.parse(buffer, apiVersion);
case CREATE_TOPICS:
- return new CreateTopicsRequest(struct, apiVersion);
+ return CreateTopicsRequest.parse(buffer, apiVersion);
case DELETE_TOPICS:
- return new DeleteTopicsRequest(struct, apiVersion);
+ return DeleteTopicsRequest.parse(buffer, apiVersion);
case DELETE_RECORDS:
- return new DeleteRecordsRequest(struct, apiVersion);
+ return DeleteRecordsRequest.parse(buffer, apiVersion);
case INIT_PRODUCER_ID:
- return new InitProducerIdRequest(struct, apiVersion);
+ return InitProducerIdRequest.parse(buffer, apiVersion);
case OFFSET_FOR_LEADER_EPOCH:
- return new OffsetsForLeaderEpochRequest(struct, apiVersion);
+ return OffsetsForLeaderEpochRequest.parse(buffer, apiVersion);
case ADD_PARTITIONS_TO_TXN:
- return new AddPartitionsToTxnRequest(struct, apiVersion);
+ return AddPartitionsToTxnRequest.parse(buffer, apiVersion);
case ADD_OFFSETS_TO_TXN:
- return new AddOffsetsToTxnRequest(struct, apiVersion);
+ return AddOffsetsToTxnRequest.parse(buffer, apiVersion);
case END_TXN:
- return new EndTxnRequest(struct, apiVersion);
+ return EndTxnRequest.parse(buffer, apiVersion);
case WRITE_TXN_MARKERS:
- return new WriteTxnMarkersRequest(struct, apiVersion);
+ return WriteTxnMarkersRequest.parse(buffer, apiVersion);
case TXN_OFFSET_COMMIT:
- return new TxnOffsetCommitRequest(struct, apiVersion);
+ return TxnOffsetCommitRequest.parse(buffer, apiVersion);
case DESCRIBE_ACLS:
- return new DescribeAclsRequest(struct, apiVersion);
+ return DescribeAclsRequest.parse(buffer, apiVersion);
case CREATE_ACLS:
- return new CreateAclsRequest(struct, apiVersion);
+ return CreateAclsRequest.parse(buffer, apiVersion);
case DELETE_ACLS:
- return new DeleteAclsRequest(struct, apiVersion);
+ return DeleteAclsRequest.parse(buffer, apiVersion);
case DESCRIBE_CONFIGS:
- return new DescribeConfigsRequest(struct, apiVersion);
+ return DescribeConfigsRequest.parse(buffer, apiVersion);
case ALTER_CONFIGS:
- return new AlterConfigsRequest(struct, apiVersion);
+ return AlterConfigsRequest.parse(buffer, apiVersion);
case ALTER_REPLICA_LOG_DIRS:
- return new AlterReplicaLogDirsRequest(struct, apiVersion);
+ return AlterReplicaLogDirsRequest.parse(buffer, apiVersion);
case DESCRIBE_LOG_DIRS:
- return new DescribeLogDirsRequest(struct, apiVersion);
+ return DescribeLogDirsRequest.parse(buffer, apiVersion);
case SASL_AUTHENTICATE:
- return new SaslAuthenticateRequest(struct, apiVersion);
+ return SaslAuthenticateRequest.parse(buffer, apiVersion);
case CREATE_PARTITIONS:
- return new CreatePartitionsRequest(struct, apiVersion);
+ return CreatePartitionsRequest.parse(buffer, apiVersion);
case CREATE_DELEGATION_TOKEN:
- return new CreateDelegationTokenRequest(struct, apiVersion);
+ return CreateDelegationTokenRequest.parse(buffer, apiVersion);
case RENEW_DELEGATION_TOKEN:
- return new RenewDelegationTokenRequest(struct, apiVersion);
+ return RenewDelegationTokenRequest.parse(buffer, apiVersion);
case EXPIRE_DELEGATION_TOKEN:
- return new ExpireDelegationTokenRequest(struct, apiVersion);
+ return ExpireDelegationTokenRequest.parse(buffer, apiVersion);
case DESCRIBE_DELEGATION_TOKEN:
- return new DescribeDelegationTokenRequest(struct, apiVersion);
+ return DescribeDelegationTokenRequest.parse(buffer, apiVersion);
case DELETE_GROUPS:
- return new DeleteGroupsRequest(struct, apiVersion);
+ return DeleteGroupsRequest.parse(buffer, apiVersion);
case ELECT_LEADERS:
- return new ElectLeadersRequest(struct, apiVersion);
+ return ElectLeadersRequest.parse(buffer, apiVersion);
case INCREMENTAL_ALTER_CONFIGS:
- return new IncrementalAlterConfigsRequest(struct, apiVersion);
+ return IncrementalAlterConfigsRequest.parse(buffer, apiVersion);
case ALTER_PARTITION_REASSIGNMENTS:
- return new AlterPartitionReassignmentsRequest(struct, apiVersion);
+ return AlterPartitionReassignmentsRequest.parse(buffer, apiVersion);
case LIST_PARTITION_REASSIGNMENTS:
- return new ListPartitionReassignmentsRequest(struct, apiVersion);
+ return ListPartitionReassignmentsRequest.parse(buffer, apiVersion);
case OFFSET_DELETE:
- return new OffsetDeleteRequest(struct, apiVersion);
+ return OffsetDeleteRequest.parse(buffer, apiVersion);
case DESCRIBE_CLIENT_QUOTAS:
- return new DescribeClientQuotasRequest(struct, apiVersion);
+ return DescribeClientQuotasRequest.parse(buffer, apiVersion);
case ALTER_CLIENT_QUOTAS:
- return new AlterClientQuotasRequest(struct, apiVersion);
+ return AlterClientQuotasRequest.parse(buffer, apiVersion);
case DESCRIBE_USER_SCRAM_CREDENTIALS:
- return new DescribeUserScramCredentialsRequest(struct, apiVersion);
+ return DescribeUserScramCredentialsRequest.parse(buffer, apiVersion);
case ALTER_USER_SCRAM_CREDENTIALS:
- return new AlterUserScramCredentialsRequest(struct, apiVersion);
+ return AlterUserScramCredentialsRequest.parse(buffer, apiVersion);
case VOTE:
- return new VoteRequest(struct, apiVersion);
+ return VoteRequest.parse(buffer, apiVersion);
case BEGIN_QUORUM_EPOCH:
- return new BeginQuorumEpochRequest(struct, apiVersion);
+ return BeginQuorumEpochRequest.parse(buffer, apiVersion);
case END_QUORUM_EPOCH:
- return new EndQuorumEpochRequest(struct, apiVersion);
+ return EndQuorumEpochRequest.parse(buffer, apiVersion);
case DESCRIBE_QUORUM:
- return new DescribeQuorumRequest(struct, apiVersion);
+ return DescribeQuorumRequest.parse(buffer, apiVersion);
case ALTER_ISR:
- return new AlterIsrRequest(new AlterIsrRequestData(struct, apiVersion), apiVersion);
+ return AlterIsrRequest.parse(buffer, apiVersion);
case UPDATE_FEATURES:
- return new UpdateFeaturesRequest(struct, apiVersion);
+ return UpdateFeaturesRequest.parse(buffer, apiVersion);
case ENVELOPE:
- return new EnvelopeRequest(struct, apiVersion);
+ return EnvelopeRequest.parse(buffer, apiVersion);
default:
throw new AssertionError(String.format("ApiKey %s is not currently handled in `parseRequest`, the " +
"code should be updated to do so.", apiKey));
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequestResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequestResponse.java
index 39698a1d20b74..5d655bd0075dd 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequestResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequestResponse.java
@@ -16,6 +16,4 @@
*/
package org.apache.kafka.common.requests;
-public interface AbstractRequestResponse {
-
-}
+public interface AbstractRequestResponse { }
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java
index 11566073374cf..747698166359c 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java
@@ -16,38 +16,49 @@
*/
package org.apache.kafka.common.requests;
-import org.apache.kafka.common.message.EnvelopeResponseData;
-import org.apache.kafka.common.message.FetchResponseData;
-import org.apache.kafka.common.message.AlterIsrResponseData;
-import org.apache.kafka.common.message.ProduceResponseData;
-import org.apache.kafka.common.network.NetworkSend;
import org.apache.kafka.common.network.Send;
import org.apache.kafka.common.protocol.ApiKeys;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.Message;
+import org.apache.kafka.common.protocol.SendBuilder;
import java.nio.ByteBuffer;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
+import java.util.Objects;
import java.util.stream.Collectors;
import java.util.stream.Stream;
public abstract class AbstractResponse implements AbstractRequestResponse {
public static final int DEFAULT_THROTTLE_TIME = 0;
- protected Send toSend(String destination, ResponseHeader header, short apiVersion) {
- return new NetworkSend(destination, RequestUtils.serialize(header.toStruct(), toStruct(apiVersion)));
+ private final ApiKeys apiKey;
+
+ protected AbstractResponse(ApiKeys apiKey) {
+ this.apiKey = apiKey;
+ }
+
+ public final Send toSend(String destination, ResponseHeader header, short version) {
+ return SendBuilder.buildResponseSend(destination, header, data(), version);
}
/**
* Visible for testing, typically {@link #toSend(String, ResponseHeader, short)} should be used instead.
*/
- public ByteBuffer serialize(ApiKeys apiKey, short version, int correlationId) {
- ResponseHeader header =
- new ResponseHeader(correlationId, apiKey.responseHeaderVersion(version));
- return RequestUtils.serialize(header.toStruct(), toStruct(version));
+ public final ByteBuffer serializeWithHeader(short version, int correlationId) {
+ return serializeWithHeader(new ResponseHeader(correlationId, apiKey.responseHeaderVersion(version)), version);
+ }
+
+ final ByteBuffer serializeWithHeader(ResponseHeader header, short version) {
+ Objects.requireNonNull(header, "header should not be null");
+ return RequestUtils.serialize(header.data(), header.headerVersion(), data(), version);
+ }
+
+ // Visible for testing
+ final ByteBuffer serializeBody(short version) {
+ return RequestUtils.serialize(null, (short) 0, data(), version);
}
/**
@@ -84,17 +95,18 @@ protected void updateErrorCounts(Map errorCounts, Errors error)
errorCounts.put(error, count + 1);
}
- protected abstract Struct toStruct(short version);
+ protected abstract Message data();
/**
* Parse a response from the provided buffer. The buffer is expected to hold both
* the {@link ResponseHeader} as well as the response payload.
*/
- public static AbstractResponse parseResponse(ByteBuffer byteBuffer, RequestHeader requestHeader) {
+ public static AbstractResponse parseResponse(ByteBuffer buffer, RequestHeader requestHeader) {
ApiKeys apiKey = requestHeader.apiKey();
short apiVersion = requestHeader.apiVersion();
- ResponseHeader responseHeader = ResponseHeader.parse(byteBuffer, apiKey.responseHeaderVersion(apiVersion));
+ ResponseHeader responseHeader = ResponseHeader.parse(buffer, apiKey.responseHeaderVersion(apiVersion));
+
if (requestHeader.correlationId() != responseHeader.correlationId()) {
throw new CorrelationIdMismatchException("Correlation id for response ("
+ responseHeader.correlationId() + ") does not match request ("
@@ -102,130 +114,129 @@ public static AbstractResponse parseResponse(ByteBuffer byteBuffer, RequestHeade
requestHeader.correlationId(), responseHeader.correlationId());
}
- Struct struct = apiKey.parseResponse(apiVersion, byteBuffer);
- return AbstractResponse.parseResponse(apiKey, struct, apiVersion);
+ return AbstractResponse.parseResponse(apiKey, buffer, apiVersion);
}
- public static AbstractResponse parseResponse(ApiKeys apiKey, Struct struct, short version) {
+ public static AbstractResponse parseResponse(ApiKeys apiKey, ByteBuffer responseBuffer, short version) {
switch (apiKey) {
case PRODUCE:
- return new ProduceResponse(new ProduceResponseData(struct, version));
+ return ProduceResponse.parse(responseBuffer, version);
case FETCH:
- return new FetchResponse<>(new FetchResponseData(struct, version));
+ return FetchResponse.parse(responseBuffer, version);
case LIST_OFFSETS:
- return new ListOffsetResponse(struct, version);
+ return ListOffsetResponse.parse(responseBuffer, version);
case METADATA:
- return new MetadataResponse(struct, version);
+ return MetadataResponse.parse(responseBuffer, version);
case OFFSET_COMMIT:
- return new OffsetCommitResponse(struct, version);
+ return OffsetCommitResponse.parse(responseBuffer, version);
case OFFSET_FETCH:
- return new OffsetFetchResponse(struct, version);
+ return OffsetFetchResponse.parse(responseBuffer, version);
case FIND_COORDINATOR:
- return new FindCoordinatorResponse(struct, version);
+ return FindCoordinatorResponse.parse(responseBuffer, version);
case JOIN_GROUP:
- return new JoinGroupResponse(struct, version);
+ return JoinGroupResponse.parse(responseBuffer, version);
case HEARTBEAT:
- return new HeartbeatResponse(struct, version);
+ return HeartbeatResponse.parse(responseBuffer, version);
case LEAVE_GROUP:
- return new LeaveGroupResponse(struct, version);
+ return LeaveGroupResponse.parse(responseBuffer, version);
case SYNC_GROUP:
- return new SyncGroupResponse(struct, version);
+ return SyncGroupResponse.parse(responseBuffer, version);
case STOP_REPLICA:
- return new StopReplicaResponse(struct, version);
+ return StopReplicaResponse.parse(responseBuffer, version);
case CONTROLLED_SHUTDOWN:
- return new ControlledShutdownResponse(struct, version);
+ return ControlledShutdownResponse.parse(responseBuffer, version);
case UPDATE_METADATA:
- return new UpdateMetadataResponse(struct, version);
+ return UpdateMetadataResponse.parse(responseBuffer, version);
case LEADER_AND_ISR:
- return new LeaderAndIsrResponse(struct, version);
+ return LeaderAndIsrResponse.parse(responseBuffer, version);
case DESCRIBE_GROUPS:
- return new DescribeGroupsResponse(struct, version);
+ return DescribeGroupsResponse.parse(responseBuffer, version);
case LIST_GROUPS:
- return new ListGroupsResponse(struct, version);
+ return ListGroupsResponse.parse(responseBuffer, version);
case SASL_HANDSHAKE:
- return new SaslHandshakeResponse(struct, version);
+ return SaslHandshakeResponse.parse(responseBuffer, version);
case API_VERSIONS:
- return ApiVersionsResponse.fromStruct(struct, version);
+ return ApiVersionsResponse.parse(responseBuffer, version);
case CREATE_TOPICS:
- return new CreateTopicsResponse(struct, version);
+ return CreateTopicsResponse.parse(responseBuffer, version);
case DELETE_TOPICS:
- return new DeleteTopicsResponse(struct, version);
+ return DeleteTopicsResponse.parse(responseBuffer, version);
case DELETE_RECORDS:
- return new DeleteRecordsResponse(struct, version);
+ return DeleteRecordsResponse.parse(responseBuffer, version);
case INIT_PRODUCER_ID:
- return new InitProducerIdResponse(struct, version);
+ return InitProducerIdResponse.parse(responseBuffer, version);
case OFFSET_FOR_LEADER_EPOCH:
- return new OffsetsForLeaderEpochResponse(struct, version);
+ return OffsetsForLeaderEpochResponse.parse(responseBuffer, version);
case ADD_PARTITIONS_TO_TXN:
- return new AddPartitionsToTxnResponse(struct, version);
+ return AddPartitionsToTxnResponse.parse(responseBuffer, version);
case ADD_OFFSETS_TO_TXN:
- return new AddOffsetsToTxnResponse(struct, version);
+ return AddOffsetsToTxnResponse.parse(responseBuffer, version);
case END_TXN:
- return new EndTxnResponse(struct, version);
+ return EndTxnResponse.parse(responseBuffer, version);
case WRITE_TXN_MARKERS:
- return new WriteTxnMarkersResponse(struct, version);
+ return WriteTxnMarkersResponse.parse(responseBuffer, version);
case TXN_OFFSET_COMMIT:
- return new TxnOffsetCommitResponse(struct, version);
+ return TxnOffsetCommitResponse.parse(responseBuffer, version);
case DESCRIBE_ACLS:
- return new DescribeAclsResponse(struct, version);
+ return DescribeAclsResponse.parse(responseBuffer, version);
case CREATE_ACLS:
- return new CreateAclsResponse(struct, version);
+ return CreateAclsResponse.parse(responseBuffer, version);
case DELETE_ACLS:
- return new DeleteAclsResponse(struct, version);
+ return DeleteAclsResponse.parse(responseBuffer, version);
case DESCRIBE_CONFIGS:
- return new DescribeConfigsResponse(struct, version);
+ return DescribeConfigsResponse.parse(responseBuffer, version);
case ALTER_CONFIGS:
- return new AlterConfigsResponse(struct, version);
+ return AlterConfigsResponse.parse(responseBuffer, version);
case ALTER_REPLICA_LOG_DIRS:
- return new AlterReplicaLogDirsResponse(struct, version);
+ return AlterReplicaLogDirsResponse.parse(responseBuffer, version);
case DESCRIBE_LOG_DIRS:
- return new DescribeLogDirsResponse(struct, version);
+ return DescribeLogDirsResponse.parse(responseBuffer, version);
case SASL_AUTHENTICATE:
- return new SaslAuthenticateResponse(struct, version);
+ return SaslAuthenticateResponse.parse(responseBuffer, version);
case CREATE_PARTITIONS:
- return new CreatePartitionsResponse(struct, version);
+ return CreatePartitionsResponse.parse(responseBuffer, version);
case CREATE_DELEGATION_TOKEN:
- return new CreateDelegationTokenResponse(struct, version);
+ return CreateDelegationTokenResponse.parse(responseBuffer, version);
case RENEW_DELEGATION_TOKEN:
- return new RenewDelegationTokenResponse(struct, version);
+ return RenewDelegationTokenResponse.parse(responseBuffer, version);
case EXPIRE_DELEGATION_TOKEN:
- return new ExpireDelegationTokenResponse(struct, version);
+ return ExpireDelegationTokenResponse.parse(responseBuffer, version);
case DESCRIBE_DELEGATION_TOKEN:
- return new DescribeDelegationTokenResponse(struct, version);
+ return DescribeDelegationTokenResponse.parse(responseBuffer, version);
case DELETE_GROUPS:
- return new DeleteGroupsResponse(struct, version);
+ return DeleteGroupsResponse.parse(responseBuffer, version);
case ELECT_LEADERS:
- return new ElectLeadersResponse(struct, version);
+ return ElectLeadersResponse.parse(responseBuffer, version);
case INCREMENTAL_ALTER_CONFIGS:
- return new IncrementalAlterConfigsResponse(struct, version);
+ return IncrementalAlterConfigsResponse.parse(responseBuffer, version);
case ALTER_PARTITION_REASSIGNMENTS:
- return new AlterPartitionReassignmentsResponse(struct, version);
+ return AlterPartitionReassignmentsResponse.parse(responseBuffer, version);
case LIST_PARTITION_REASSIGNMENTS:
- return new ListPartitionReassignmentsResponse(struct, version);
+ return ListPartitionReassignmentsResponse.parse(responseBuffer, version);
case OFFSET_DELETE:
- return new OffsetDeleteResponse(struct, version);
+ return OffsetDeleteResponse.parse(responseBuffer, version);
case DESCRIBE_CLIENT_QUOTAS:
- return new DescribeClientQuotasResponse(struct, version);
+ return DescribeClientQuotasResponse.parse(responseBuffer, version);
case ALTER_CLIENT_QUOTAS:
- return new AlterClientQuotasResponse(struct, version);
+ return AlterClientQuotasResponse.parse(responseBuffer, version);
case DESCRIBE_USER_SCRAM_CREDENTIALS:
- return new DescribeUserScramCredentialsResponse(struct, version);
+ return DescribeUserScramCredentialsResponse.parse(responseBuffer, version);
case ALTER_USER_SCRAM_CREDENTIALS:
- return new AlterUserScramCredentialsResponse(struct, version);
+ return AlterUserScramCredentialsResponse.parse(responseBuffer, version);
case VOTE:
- return new VoteResponse(struct, version);
+ return VoteResponse.parse(responseBuffer, version);
case BEGIN_QUORUM_EPOCH:
- return new BeginQuorumEpochResponse(struct, version);
+ return BeginQuorumEpochResponse.parse(responseBuffer, version);
case END_QUORUM_EPOCH:
- return new EndQuorumEpochResponse(struct, version);
+ return EndQuorumEpochResponse.parse(responseBuffer, version);
case DESCRIBE_QUORUM:
- return new DescribeQuorumResponse(struct, version);
+ return DescribeQuorumResponse.parse(responseBuffer, version);
case ALTER_ISR:
- return new AlterIsrResponse(new AlterIsrResponseData(struct, version));
+ return AlterIsrResponse.parse(responseBuffer, version);
case UPDATE_FEATURES:
- return new UpdateFeaturesResponse(struct, version);
+ return UpdateFeaturesResponse.parse(responseBuffer, version);
case ENVELOPE:
- return new EnvelopeResponse(new EnvelopeResponseData(struct, version));
+ return EnvelopeResponse.parse(responseBuffer, version);
default:
throw new AssertionError(String.format("ApiKey %s is not currently handled in `parseResponse`, the " +
"code should be updated to do so.", apiKey));
@@ -241,11 +252,13 @@ public boolean shouldClientThrottle(short version) {
return false;
}
- public int throttleTimeMs() {
- return DEFAULT_THROTTLE_TIME;
+ public ApiKeys apiKey() {
+ return apiKey;
}
- public String toString(short version) {
- return toStruct(version).toString();
+ public abstract int throttleTimeMs();
+
+ public String toString() {
+ return data().toString();
}
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AddOffsetsToTxnRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/AddOffsetsToTxnRequest.java
index 3b1c746c5b5c9..4ad56aa94268f 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AddOffsetsToTxnRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AddOffsetsToTxnRequest.java
@@ -19,8 +19,8 @@
import org.apache.kafka.common.message.AddOffsetsToTxnRequestData;
import org.apache.kafka.common.message.AddOffsetsToTxnResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
@@ -52,14 +52,9 @@ public AddOffsetsToTxnRequest(AddOffsetsToTxnRequestData data, short version) {
this.data = data;
}
- public AddOffsetsToTxnRequest(Struct struct, short version) {
- super(ApiKeys.ADD_OFFSETS_TO_TXN, version);
- this.data = new AddOffsetsToTxnRequestData(struct, version);
- }
-
@Override
- protected Struct toStruct() {
- return data.toStruct(version());
+ protected AddOffsetsToTxnRequestData data() {
+ return data;
}
@Override
@@ -70,6 +65,6 @@ public AddOffsetsToTxnResponse getErrorResponse(int throttleTimeMs, Throwable e)
}
public static AddOffsetsToTxnRequest parse(ByteBuffer buffer, short version) {
- return new AddOffsetsToTxnRequest(ApiKeys.ADD_OFFSETS_TO_TXN.parseRequest(version, buffer), version);
+ return new AddOffsetsToTxnRequest(new AddOffsetsToTxnRequestData(new ByteBufferAccessor(buffer), version), version);
}
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AddOffsetsToTxnResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AddOffsetsToTxnResponse.java
index 99864f4bffd4a..6908e51960609 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AddOffsetsToTxnResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AddOffsetsToTxnResponse.java
@@ -18,8 +18,8 @@
import org.apache.kafka.common.message.AddOffsetsToTxnResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.Map;
@@ -42,30 +42,27 @@ public class AddOffsetsToTxnResponse extends AbstractResponse {
public AddOffsetsToTxnResponseData data;
public AddOffsetsToTxnResponse(AddOffsetsToTxnResponseData data) {
+ super(ApiKeys.ADD_OFFSETS_TO_TXN);
this.data = data;
}
- public AddOffsetsToTxnResponse(Struct struct, short version) {
- this.data = new AddOffsetsToTxnResponseData(struct, version);
- }
-
@Override
public Map errorCounts() {
return errorCounts(Errors.forCode(data.errorCode()));
}
@Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
+ public int throttleTimeMs() {
+ return data.throttleTimeMs();
}
@Override
- public int throttleTimeMs() {
- return data.throttleTimeMs();
+ protected AddOffsetsToTxnResponseData data() {
+ return data;
}
public static AddOffsetsToTxnResponse parse(ByteBuffer buffer, short version) {
- return new AddOffsetsToTxnResponse(ApiKeys.ADD_OFFSETS_TO_TXN.parseResponse(version, buffer), version);
+ return new AddOffsetsToTxnResponse(new AddOffsetsToTxnResponseData(new ByteBufferAccessor(buffer), version));
}
@Override
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AddPartitionsToTxnRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/AddPartitionsToTxnRequest.java
index 13be2127cdd4d..57645789ca000 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AddPartitionsToTxnRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AddPartitionsToTxnRequest.java
@@ -21,8 +21,8 @@
import org.apache.kafka.common.message.AddPartitionsToTxnRequestData.AddPartitionsToTxnTopic;
import org.apache.kafka.common.message.AddPartitionsToTxnRequestData.AddPartitionsToTxnTopicCollection;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.ArrayList;
@@ -103,11 +103,6 @@ public AddPartitionsToTxnRequest(final AddPartitionsToTxnRequestData data, short
this.data = data;
}
- public AddPartitionsToTxnRequest(Struct struct, short version) {
- super(ApiKeys.ADD_PARTITIONS_TO_TXN, version);
- this.data = new AddPartitionsToTxnRequestData(struct, version);
- }
-
public List partitions() {
if (cachedPartitions != null) {
return cachedPartitions;
@@ -117,8 +112,8 @@ public List partitions() {
}
@Override
- protected Struct toStruct() {
- return data.toStruct(version());
+ protected AddPartitionsToTxnRequestData data() {
+ return data;
}
@Override
@@ -131,6 +126,6 @@ public AddPartitionsToTxnResponse getErrorResponse(int throttleTimeMs, Throwable
}
public static AddPartitionsToTxnRequest parse(ByteBuffer buffer, short version) {
- return new AddPartitionsToTxnRequest(ApiKeys.ADD_PARTITIONS_TO_TXN.parseRequest(version, buffer), version);
+ return new AddPartitionsToTxnRequest(new AddPartitionsToTxnRequestData(new ByteBufferAccessor(buffer), version), version);
}
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AddPartitionsToTxnResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AddPartitionsToTxnResponse.java
index cc79ce8ca4200..c6da3af0a9319 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AddPartitionsToTxnResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AddPartitionsToTxnResponse.java
@@ -23,8 +23,8 @@
import org.apache.kafka.common.message.AddPartitionsToTxnResponseData.AddPartitionsToTxnTopicResult;
import org.apache.kafka.common.message.AddPartitionsToTxnResponseData.AddPartitionsToTxnTopicResultCollection;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.HashMap;
@@ -50,11 +50,13 @@ public class AddPartitionsToTxnResponse extends AbstractResponse {
private Map cachedErrorsMap = null;
- public AddPartitionsToTxnResponse(Struct struct, short version) {
- this.data = new AddPartitionsToTxnResponseData(struct, version);
+ public AddPartitionsToTxnResponse(AddPartitionsToTxnResponseData data) {
+ super(ApiKeys.ADD_PARTITIONS_TO_TXN);
+ this.data = data;
}
public AddPartitionsToTxnResponse(int throttleTimeMs, Map errors) {
+ super(ApiKeys.ADD_PARTITIONS_TO_TXN);
Map resultMap = new HashMap<>();
@@ -115,12 +117,12 @@ public Map errorCounts() {
}
@Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
+ protected AddPartitionsToTxnResponseData data() {
+ return data;
}
public static AddPartitionsToTxnResponse parse(ByteBuffer buffer, short version) {
- return new AddPartitionsToTxnResponse(ApiKeys.ADD_PARTITIONS_TO_TXN.parseResponse(version, buffer), version);
+ return new AddPartitionsToTxnResponse(new AddPartitionsToTxnResponseData(new ByteBufferAccessor(buffer), version));
}
@Override
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterClientQuotasRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterClientQuotasRequest.java
index 6d4c2a1f357c3..7178475123633 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AlterClientQuotasRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterClientQuotasRequest.java
@@ -20,11 +20,13 @@
import org.apache.kafka.common.message.AlterClientQuotasRequestData.EntityData;
import org.apache.kafka.common.message.AlterClientQuotasRequestData.EntryData;
import org.apache.kafka.common.message.AlterClientQuotasRequestData.OpData;
+import org.apache.kafka.common.message.AlterClientQuotasResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.quota.ClientQuotaAlteration;
import org.apache.kafka.common.quota.ClientQuotaEntity;
+import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Collection;
import java.util.HashMap;
@@ -85,12 +87,7 @@ public AlterClientQuotasRequest(AlterClientQuotasRequestData data, short version
this.data = data;
}
- public AlterClientQuotasRequest(Struct struct, short version) {
- super(ApiKeys.ALTER_CLIENT_QUOTAS, version);
- this.data = new AlterClientQuotasRequestData(struct, version);
- }
-
- public Collection entries() {
+ public List entries() {
List entries = new ArrayList<>(data.entries().size());
for (EntryData entryData : data.entries()) {
Map entity = new HashMap<>(entryData.entity().size());
@@ -113,21 +110,30 @@ public boolean validateOnly() {
return data.validateOnly();
}
+ @Override
+ protected AlterClientQuotasRequestData data() {
+ return data;
+ }
+
@Override
public AlterClientQuotasResponse getErrorResponse(int throttleTimeMs, Throwable e) {
- ArrayList entities = new ArrayList<>(data.entries().size());
+ List responseEntries = new ArrayList<>();
for (EntryData entryData : data.entries()) {
- Map entity = new HashMap<>(entryData.entity().size());
+ List responseEntities = new ArrayList<>();
for (EntityData entityData : entryData.entity()) {
- entity.put(entityData.entityType(), entityData.entityName());
+ responseEntities.add(new AlterClientQuotasResponseData.EntityData()
+ .setEntityType(entityData.entityType())
+ .setEntityName(entityData.entityName()));
}
- entities.add(new ClientQuotaEntity(entity));
+ responseEntries.add(new AlterClientQuotasResponseData.EntryData().setEntity(responseEntities));
}
- return new AlterClientQuotasResponse(entities, throttleTimeMs, e);
+ AlterClientQuotasResponseData responseData = new AlterClientQuotasResponseData()
+ .setThrottleTimeMs(throttleTimeMs)
+ .setEntries(responseEntries);
+ return new AlterClientQuotasResponse(responseData);
}
- @Override
- protected Struct toStruct() {
- return data.toStruct(version());
+ public static AlterClientQuotasRequest parse(ByteBuffer buffer, short version) {
+ return new AlterClientQuotasRequest(new AlterClientQuotasRequestData(new ByteBufferAccessor(buffer), version), version);
}
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterClientQuotasResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterClientQuotasResponse.java
index 3e95671a6f939..bce548822d895 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AlterClientQuotasResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterClientQuotasResponse.java
@@ -21,13 +21,12 @@
import org.apache.kafka.common.message.AlterClientQuotasResponseData.EntityData;
import org.apache.kafka.common.message.AlterClientQuotasResponseData.EntryData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import org.apache.kafka.common.quota.ClientQuotaEntity;
import java.nio.ByteBuffer;
import java.util.ArrayList;
-import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -36,39 +35,9 @@ public class AlterClientQuotasResponse extends AbstractResponse {
private final AlterClientQuotasResponseData data;
- public AlterClientQuotasResponse(Map result, int throttleTimeMs) {
- List entries = new ArrayList<>(result.size());
- for (Map.Entry entry : result.entrySet()) {
- ApiError e = entry.getValue();
- entries.add(new EntryData()
- .setErrorCode(e.error().code())
- .setErrorMessage(e.message())
- .setEntity(toEntityData(entry.getKey())));
- }
-
- this.data = new AlterClientQuotasResponseData()
- .setThrottleTimeMs(throttleTimeMs)
- .setEntries(entries);
- }
-
- public AlterClientQuotasResponse(Collection entities, int throttleTimeMs, Throwable e) {
- ApiError apiError = ApiError.fromThrowable(e);
-
- List entries = new ArrayList<>(entities.size());
- for (ClientQuotaEntity entity : entities) {
- entries.add(new EntryData()
- .setErrorCode(apiError.error().code())
- .setErrorMessage(apiError.message())
- .setEntity(toEntityData(entity)));
- }
-
- this.data = new AlterClientQuotasResponseData()
- .setThrottleTimeMs(throttleTimeMs)
- .setEntries(entries);
- }
-
- public AlterClientQuotasResponse(Struct struct, short version) {
- this.data = new AlterClientQuotasResponseData(struct, version);
+ public AlterClientQuotasResponse(AlterClientQuotasResponseData data) {
+ super(ApiKeys.ALTER_CLIENT_QUOTAS);
+ this.data = data;
}
public void complete(Map> futures) {
@@ -108,8 +77,8 @@ public Map errorCounts() {
}
@Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
+ protected AlterClientQuotasResponseData data() {
+ return data;
}
private static List toEntityData(ClientQuotaEntity entity) {
@@ -123,6 +92,22 @@ private static List toEntityData(ClientQuotaEntity entity) {
}
public static AlterClientQuotasResponse parse(ByteBuffer buffer, short version) {
- return new AlterClientQuotasResponse(ApiKeys.ALTER_CLIENT_QUOTAS.parseResponse(version, buffer), version);
+ return new AlterClientQuotasResponse(new AlterClientQuotasResponseData(new ByteBufferAccessor(buffer), version));
}
+
+ public static AlterClientQuotasResponse fromQuotaEntities(Map result, int throttleTimeMs) {
+ List entries = new ArrayList<>(result.size());
+ for (Map.Entry entry : result.entrySet()) {
+ ApiError e = entry.getValue();
+ entries.add(new EntryData()
+ .setErrorCode(e.error().code())
+ .setErrorMessage(e.message())
+ .setEntity(toEntityData(entry.getKey())));
+ }
+
+ return new AlterClientQuotasResponse(new AlterClientQuotasResponseData()
+ .setThrottleTimeMs(throttleTimeMs)
+ .setEntries(entries));
+ }
+
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterConfigsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterConfigsRequest.java
index 1ee3330a71ee5..f30e8b9c1ffe9 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AlterConfigsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterConfigsRequest.java
@@ -21,7 +21,7 @@
import org.apache.kafka.common.message.AlterConfigsRequestData;
import org.apache.kafka.common.message.AlterConfigsResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import java.nio.ByteBuffer;
import java.util.Collection;
@@ -97,11 +97,6 @@ public AlterConfigsRequest(AlterConfigsRequestData data, short version) {
this.data = data;
}
- public AlterConfigsRequest(Struct struct, short version) {
- super(ApiKeys.ALTER_CONFIGS, version);
- this.data = new AlterConfigsRequestData(struct, version);
- }
-
public Map configs() {
return data.resources().stream().collect(Collectors.toMap(
resource -> new ConfigResource(
@@ -117,8 +112,8 @@ public boolean validateOnly() {
}
@Override
- protected Struct toStruct() {
- return data.toStruct(version());
+ protected AlterConfigsRequestData data() {
+ return data;
}
@Override
@@ -138,6 +133,6 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
}
public static AlterConfigsRequest parse(ByteBuffer buffer, short version) {
- return new AlterConfigsRequest(ApiKeys.ALTER_CONFIGS.parseRequest(version, buffer), version);
+ return new AlterConfigsRequest(new AlterConfigsRequestData(new ByteBufferAccessor(buffer), version), version);
}
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterConfigsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterConfigsResponse.java
index 7122a132e5393..1115f06ee80a9 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AlterConfigsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterConfigsResponse.java
@@ -20,8 +20,8 @@
import org.apache.kafka.common.config.ConfigResource;
import org.apache.kafka.common.message.AlterConfigsResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.Map;
@@ -32,17 +32,10 @@ public class AlterConfigsResponse extends AbstractResponse {
private final AlterConfigsResponseData data;
public AlterConfigsResponse(AlterConfigsResponseData data) {
+ super(ApiKeys.ALTER_CONFIGS);
this.data = data;
}
- public AlterConfigsResponse(Struct struct, short version) {
- this.data = new AlterConfigsResponseData(struct, version);
- }
-
- public AlterConfigsResponseData data() {
- return data;
- }
-
public Map errors() {
return data.responses().stream().collect(Collectors.toMap(
response -> new ConfigResource(
@@ -63,12 +56,12 @@ public int throttleTimeMs() {
}
@Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
+ public AlterConfigsResponseData data() {
+ return data;
}
public static AlterConfigsResponse parse(ByteBuffer buffer, short version) {
- return new AlterConfigsResponse(ApiKeys.ALTER_CONFIGS.parseResponse(version, buffer), version);
+ return new AlterConfigsResponse(new AlterConfigsResponseData(new ByteBufferAccessor(buffer), version));
}
@Override
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterIsrRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterIsrRequest.java
index 61662721f19b6..7ce86aedf432e 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AlterIsrRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterIsrRequest.java
@@ -20,8 +20,10 @@
import org.apache.kafka.common.message.AlterIsrRequestData;
import org.apache.kafka.common.message.AlterIsrResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
+
+import java.nio.ByteBuffer;
public class AlterIsrRequest extends AbstractRequest {
@@ -36,11 +38,6 @@ public AlterIsrRequestData data() {
return data;
}
- @Override
- protected Struct toStruct() {
- return data.toStruct(version());
- }
-
/**
* Get an error response for a request with specified throttle time in the response if applicable
*/
@@ -51,6 +48,10 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
.setErrorCode(Errors.forException(e).code()));
}
+ public static AlterIsrRequest parse(ByteBuffer buffer, short version) {
+ return new AlterIsrRequest(new AlterIsrRequestData(new ByteBufferAccessor(buffer), version), version);
+ }
+
public static class Builder extends AbstractRequest.Builder {
private final AlterIsrRequestData data;
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterIsrResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterIsrResponse.java
index b475bfaecf705..433ba66bdadd3 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AlterIsrResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterIsrResponse.java
@@ -18,9 +18,11 @@
package org.apache.kafka.common.requests;
import org.apache.kafka.common.message.AlterIsrResponseData;
+import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
+import java.nio.ByteBuffer;
import java.util.HashMap;
import java.util.Map;
@@ -29,6 +31,7 @@ public class AlterIsrResponse extends AbstractResponse {
private final AlterIsrResponseData data;
public AlterIsrResponse(AlterIsrResponseData data) {
+ super(ApiKeys.ALTER_ISR);
this.data = data;
}
@@ -46,13 +49,12 @@ public Map errorCounts() {
return counts;
}
- @Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
- }
-
@Override
public int throttleTimeMs() {
return data.throttleTimeMs();
}
+
+ public static AlterIsrResponse parse(ByteBuffer buffer, short version) {
+ return new AlterIsrResponse(new AlterIsrResponseData(new ByteBufferAccessor(buffer), version));
+ }
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterPartitionReassignmentsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterPartitionReassignmentsRequest.java
index 7b2f848f614f9..2d289cc1497e1 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AlterPartitionReassignmentsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterPartitionReassignmentsRequest.java
@@ -23,7 +23,7 @@
import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.ReassignablePartitionResponse;
import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.ReassignableTopicResponse;
import org.apache.kafka.common.protocol.ApiKeys;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import java.nio.ByteBuffer;
import java.util.ArrayList;
@@ -52,39 +52,21 @@ public String toString() {
}
private final AlterPartitionReassignmentsRequestData data;
- private final short version;
private AlterPartitionReassignmentsRequest(AlterPartitionReassignmentsRequestData data, short version) {
super(ApiKeys.ALTER_PARTITION_REASSIGNMENTS, version);
this.data = data;
- this.version = version;
- }
-
- AlterPartitionReassignmentsRequest(Struct struct, short version) {
- super(ApiKeys.ALTER_PARTITION_REASSIGNMENTS, version);
- this.data = new AlterPartitionReassignmentsRequestData(struct, version);
- this.version = version;
}
public static AlterPartitionReassignmentsRequest parse(ByteBuffer buffer, short version) {
- return new AlterPartitionReassignmentsRequest(
- ApiKeys.ALTER_PARTITION_REASSIGNMENTS.parseRequest(version, buffer),
- version
- );
+ return new AlterPartitionReassignmentsRequest(new AlterPartitionReassignmentsRequestData(
+ new ByteBufferAccessor(buffer), version), version);
}
public AlterPartitionReassignmentsRequestData data() {
return data;
}
- /**
- * Visible for testing.
- */
- @Override
- public Struct toStruct() {
- return data.toStruct(version);
- }
-
@Override
public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
ApiError apiError = ApiError.fromThrowable(e);
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterPartitionReassignmentsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterPartitionReassignmentsResponse.java
index 495258655e9bb..6aea8b1edcc87 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AlterPartitionReassignmentsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterPartitionReassignmentsResponse.java
@@ -19,8 +19,8 @@
import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.HashMap;
@@ -30,20 +30,14 @@ public class AlterPartitionReassignmentsResponse extends AbstractResponse {
private final AlterPartitionReassignmentsResponseData data;
- public AlterPartitionReassignmentsResponse(Struct struct) {
- this(struct, ApiKeys.ALTER_PARTITION_REASSIGNMENTS.latestVersion());
- }
-
public AlterPartitionReassignmentsResponse(AlterPartitionReassignmentsResponseData data) {
+ super(ApiKeys.ALTER_PARTITION_REASSIGNMENTS);
this.data = data;
}
- AlterPartitionReassignmentsResponse(Struct struct, short version) {
- this.data = new AlterPartitionReassignmentsResponseData(struct, version);
- }
-
public static AlterPartitionReassignmentsResponse parse(ByteBuffer buffer, short version) {
- return new AlterPartitionReassignmentsResponse(ApiKeys.ALTER_PARTITION_REASSIGNMENTS.responseSchema(version).read(buffer), version);
+ return new AlterPartitionReassignmentsResponse(
+ new AlterPartitionReassignmentsResponseData(new ByteBufferAccessor(buffer), version));
}
public AlterPartitionReassignmentsResponseData data() {
@@ -71,9 +65,4 @@ public Map errorCounts() {
));
return counts;
}
-
- @Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
- }
}
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 8cde2c0a1bde2..5eba039fa556b 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,18 +17,20 @@
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.ByteBufferAccessor;
+import org.apache.kafka.common.protocol.Errors;
+
import java.nio.ByteBuffer;
import java.util.HashMap;
import java.util.Map;
import java.util.stream.Collectors;
-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 {
@@ -53,22 +55,16 @@ public String toString() {
}
}
- public AlterReplicaLogDirsRequest(Struct struct, short version) {
- super(ApiKeys.ALTER_REPLICA_LOG_DIRS, version);
- this.data = new AlterReplicaLogDirsRequestData(struct, version);
- }
-
public AlterReplicaLogDirsRequest(AlterReplicaLogDirsRequestData data, short version) {
super(ApiKeys.ALTER_REPLICA_LOG_DIRS, version);
this.data = data;
}
@Override
- protected Struct toStruct() {
- return data.toStruct(version());
+ protected AlterReplicaLogDirsRequestData data() {
+ return data;
}
- @Override
public AlterReplicaLogDirsResponse getErrorResponse(int throttleTimeMs, Throwable e) {
AlterReplicaLogDirsResponseData data = new AlterReplicaLogDirsResponseData();
data.setResults(this.data.dirs().stream().flatMap(alterDir ->
@@ -93,6 +89,6 @@ public Map partitionDirs() {
}
public static AlterReplicaLogDirsRequest parse(ByteBuffer buffer, short version) {
- return new AlterReplicaLogDirsRequest(ApiKeys.ALTER_REPLICA_LOG_DIRS.parseRequest(version, buffer), version);
+ return new AlterReplicaLogDirsRequest(new AlterReplicaLogDirsRequestData(new ByteBufferAccessor(buffer), version), 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 bdc1ed0003b9e..afa658d1e150f 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,14 +17,14 @@
package org.apache.kafka.common.requests;
-import java.nio.ByteBuffer;
-import java.util.HashMap;
-import java.util.Map;
-
import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
+
+import java.nio.ByteBuffer;
+import java.util.HashMap;
+import java.util.Map;
/**
* Possible error codes:
@@ -38,23 +38,16 @@ public class AlterReplicaLogDirsResponse extends AbstractResponse {
private final AlterReplicaLogDirsResponseData data;
- public AlterReplicaLogDirsResponse(Struct struct, short version) {
- this.data = new AlterReplicaLogDirsResponseData(struct, version);
- }
-
public AlterReplicaLogDirsResponse(AlterReplicaLogDirsResponseData data) {
+ super(ApiKeys.ALTER_REPLICA_LOG_DIRS);
this.data = data;
}
+ @Override
public AlterReplicaLogDirsResponseData data() {
return data;
}
- @Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
- }
-
@Override
public int throttleTimeMs() {
return data.throttleTimeMs();
@@ -70,7 +63,7 @@ public Map errorCounts() {
}
public static AlterReplicaLogDirsResponse parse(ByteBuffer buffer, short version) {
- return new AlterReplicaLogDirsResponse(ApiKeys.ALTER_REPLICA_LOG_DIRS.responseSchema(version).read(buffer), version);
+ return new AlterReplicaLogDirsResponse(new AlterReplicaLogDirsResponseData(new ByteBufferAccessor(buffer), version));
}
@Override
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterUserScramCredentialsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterUserScramCredentialsRequest.java
index 8d3a18846021c..c319ec344525a 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AlterUserScramCredentialsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterUserScramCredentialsRequest.java
@@ -19,7 +19,7 @@
import org.apache.kafka.common.message.AlterUserScramCredentialsRequestData;
import org.apache.kafka.common.message.AlterUserScramCredentialsResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import java.nio.ByteBuffer;
import java.util.List;
@@ -48,39 +48,21 @@ public String toString() {
}
}
- private AlterUserScramCredentialsRequestData data;
- private final short version;
+ private final AlterUserScramCredentialsRequestData data;
private AlterUserScramCredentialsRequest(AlterUserScramCredentialsRequestData data, short version) {
super(ApiKeys.ALTER_USER_SCRAM_CREDENTIALS, version);
this.data = data;
- this.version = version;
- }
-
- AlterUserScramCredentialsRequest(Struct struct, short version) {
- super(ApiKeys.ALTER_USER_SCRAM_CREDENTIALS, version);
- this.data = new AlterUserScramCredentialsRequestData(struct, version);
- this.version = version;
}
public static AlterUserScramCredentialsRequest parse(ByteBuffer buffer, short version) {
- return new AlterUserScramCredentialsRequest(
- ApiKeys.ALTER_USER_SCRAM_CREDENTIALS.parseRequest(version, buffer), version
- );
+ return new AlterUserScramCredentialsRequest(new AlterUserScramCredentialsRequestData(new ByteBufferAccessor(buffer), version), version);
}
public AlterUserScramCredentialsRequestData data() {
return data;
}
- /**
- * Visible for testing.
- */
- @Override
- public Struct toStruct() {
- return data.toStruct(version);
- }
-
@Override
public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
ApiError apiError = ApiError.fromThrowable(e);
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterUserScramCredentialsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterUserScramCredentialsResponse.java
index 88ff920d3f132..2fa4937242742 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/AlterUserScramCredentialsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterUserScramCredentialsResponse.java
@@ -18,8 +18,8 @@
import org.apache.kafka.common.message.AlterUserScramCredentialsResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.Map;
@@ -28,22 +28,11 @@ public class AlterUserScramCredentialsResponse extends AbstractResponse {
private final AlterUserScramCredentialsResponseData data;
- public AlterUserScramCredentialsResponse(Struct struct) {
- this(struct, ApiKeys.ALTER_USER_SCRAM_CREDENTIALS.latestVersion());
- }
-
public AlterUserScramCredentialsResponse(AlterUserScramCredentialsResponseData responseData) {
+ super(ApiKeys.ALTER_USER_SCRAM_CREDENTIALS);
this.data = responseData;
}
- AlterUserScramCredentialsResponse(Struct struct, short version) {
- this.data = new AlterUserScramCredentialsResponseData(struct, version);
- }
-
- public static AlterUserScramCredentialsResponse parse(ByteBuffer buffer, short version) {
- return new AlterUserScramCredentialsResponse(ApiKeys.ALTER_USER_SCRAM_CREDENTIALS.responseSchema(version).read(buffer), version);
- }
-
public AlterUserScramCredentialsResponseData data() {
return data;
}
@@ -63,8 +52,7 @@ public Map errorCounts() {
return errorCounts(data.results().stream().map(r -> Errors.forCode(r.errorCode())));
}
- @Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
+ public static AlterUserScramCredentialsResponse parse(ByteBuffer buffer, short version) {
+ return new AlterUserScramCredentialsResponse(new AlterUserScramCredentialsResponseData(new ByteBufferAccessor(buffer), version));
}
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ApiError.java b/clients/src/main/java/org/apache/kafka/common/requests/ApiError.java
index a790bc8b161c6..38712adbd55f8 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/ApiError.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/ApiError.java
@@ -19,13 +19,9 @@
import org.apache.kafka.common.errors.ApiException;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.util.Objects;
-import static org.apache.kafka.common.protocol.CommonFields.ERROR_CODE;
-import static org.apache.kafka.common.protocol.CommonFields.ERROR_MESSAGE;
-
/**
* Encapsulates an error code (via the Errors enum) and an optional message. Generally, the optional message is only
* defined if it adds information over the default message associated with the error code.
@@ -47,12 +43,6 @@ public static ApiError fromThrowable(Throwable t) {
return new ApiError(error, message);
}
- public ApiError(Struct struct) {
- error = Errors.forCode(struct.get(ERROR_CODE));
- // In some cases, the error message field was introduced in newer version
- message = struct.getOrElse(ERROR_MESSAGE, null);
- }
-
public ApiError(Errors error) {
this(error, error.message());
}
@@ -67,12 +57,6 @@ public ApiError(short code, String message) {
this.message = message;
}
- public void write(Struct struct) {
- struct.set(ERROR_CODE, error.code());
- if (error != Errors.NONE)
- struct.setIfExists(ERROR_MESSAGE, message);
- }
-
public boolean is(Errors error) {
return this.error == error;
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ApiVersionsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/ApiVersionsRequest.java
index 87ffceaa26e1c..946a9d9da7bed 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/ApiVersionsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/ApiVersionsRequest.java
@@ -22,8 +22,8 @@
import org.apache.kafka.common.message.ApiVersionsResponseData.ApiVersionsResponseKey;
import org.apache.kafka.common.message.ApiVersionsResponseData.ApiVersionsResponseKeyCollection;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import org.apache.kafka.common.utils.AppInfoParser;
import java.nio.ByteBuffer;
@@ -60,7 +60,7 @@ public String toString() {
private final Short unsupportedRequestVersion;
- public final ApiVersionsRequestData data;
+ private final ApiVersionsRequestData data;
public ApiVersionsRequest(ApiVersionsRequestData data, short version) {
this(data, version, null);
@@ -78,10 +78,6 @@ public ApiVersionsRequest(ApiVersionsRequestData data, short version, Short unsu
this.unsupportedRequestVersion = unsupportedRequestVersion;
}
- public ApiVersionsRequest(Struct struct, short version) {
- this(new ApiVersionsRequestData(struct, version), version);
- }
-
public boolean hasUnsupportedRequestVersion() {
return unsupportedRequestVersion != null;
}
@@ -96,8 +92,8 @@ public boolean isValid() {
}
@Override
- protected Struct toStruct() {
- return data.toStruct(version());
+ public ApiVersionsRequestData data() {
+ return data;
}
@Override
@@ -124,7 +120,7 @@ public ApiVersionsResponse getErrorResponse(int throttleTimeMs, Throwable e) {
}
public static ApiVersionsRequest parse(ByteBuffer buffer, short version) {
- return new ApiVersionsRequest(ApiKeys.API_VERSIONS.parseRequest(version, buffer), version);
+ return new ApiVersionsRequest(new ApiVersionsRequestData(new ByteBufferAccessor(buffer), version), version);
}
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ApiVersionsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/ApiVersionsResponse.java
index 1e4ad17b47dc0..eaf8113fcd0b7 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/ApiVersionsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/ApiVersionsResponse.java
@@ -29,8 +29,6 @@
import org.apache.kafka.common.protocol.ApiKeys;
import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.SchemaException;
-import org.apache.kafka.common.protocol.types.Struct;
import org.apache.kafka.common.record.RecordBatch;
import java.nio.ByteBuffer;
@@ -45,38 +43,21 @@ public class ApiVersionsResponse extends AbstractResponse {
public static final long UNKNOWN_FINALIZED_FEATURES_EPOCH = -1L;
- public static final ApiVersionsResponse DEFAULT_API_VERSIONS_RESPONSE =
- createApiVersionsResponse(
- DEFAULT_THROTTLE_TIME,
- RecordBatch.CURRENT_MAGIC_VALUE,
- Features.emptySupportedFeatures(),
- Features.emptyFinalizedFeatures(),
- UNKNOWN_FINALIZED_FEATURES_EPOCH
- );
+ public static final ApiVersionsResponse DEFAULT_API_VERSIONS_RESPONSE = createApiVersionsResponse(
+ DEFAULT_THROTTLE_TIME, RecordBatch.CURRENT_MAGIC_VALUE);
public final ApiVersionsResponseData data;
public ApiVersionsResponse(ApiVersionsResponseData data) {
+ super(ApiKeys.API_VERSIONS);
this.data = data;
}
- public ApiVersionsResponse(Struct struct) {
- this(new ApiVersionsResponseData(struct, (short) (ApiVersionsResponseData.SCHEMAS.length - 1)));
- }
-
- public ApiVersionsResponse(Struct struct, short version) {
- this(new ApiVersionsResponseData(struct, version));
- }
-
+ @Override
public ApiVersionsResponseData data() {
return data;
}
- @Override
- protected Struct toStruct(short version) {
- return this.data.toStruct(version);
- }
-
public ApiVersionsResponseKey apiVersion(short apiKey) {
return data.apiKeys().find(apiKey);
}
@@ -100,48 +81,23 @@ public static ApiVersionsResponse parse(ByteBuffer buffer, short version) {
// Fallback to version 0 for ApiVersions response. If a client sends an ApiVersionsRequest
// using a version higher than that supported by the broker, a version 0 response is sent
// to the client indicating UNSUPPORTED_VERSION. When the client receives the response, it
- // falls back while parsing it into a Struct which means that the version received by this
- // method is not necessary the real one. It may be version 0 as well.
+ // falls back while parsing it which means that the version received by this
+ // method is not necessarily the real one. It may be version 0 as well.
int prev = buffer.position();
try {
- return new ApiVersionsResponse(
- new ApiVersionsResponseData(new ByteBufferAccessor(buffer), version));
+ return new ApiVersionsResponse(new ApiVersionsResponseData(new ByteBufferAccessor(buffer), version));
} catch (RuntimeException e) {
buffer.position(prev);
if (version != 0)
- return new ApiVersionsResponse(
- new ApiVersionsResponseData(new ByteBufferAccessor(buffer), (short) 0));
+ return new ApiVersionsResponse(new ApiVersionsResponseData(new ByteBufferAccessor(buffer), (short) 0));
else
throw e;
}
}
- public static ApiVersionsResponse fromStruct(Struct struct, short version) {
- // Fallback to version 0 for ApiVersions response. If a client sends an ApiVersionsRequest
- // using a version higher than that supported by the broker, a version 0 response is sent
- // to the client indicating UNSUPPORTED_VERSION. When the client receives the response, it
- // falls back while parsing it into a Struct which means that the version received by this
- // method is not necessary the real one. It may be version 0 as well.
- try {
- return new ApiVersionsResponse(struct, version);
- } catch (SchemaException e) {
- if (version != 0)
- return new ApiVersionsResponse(struct, (short) 0);
- else
- throw e;
- }
- }
-
- public static ApiVersionsResponse createApiVersionsResponse(
- final int throttleTimeMs,
- final byte minMagic) {
- return createApiVersionsResponse(
- throttleTimeMs,
- minMagic,
- Features.emptySupportedFeatures(),
- Features.emptyFinalizedFeatures(),
- UNKNOWN_FINALIZED_FEATURES_EPOCH
- );
+ public static ApiVersionsResponse createApiVersionsResponse(final int throttleTimeMs, final byte minMagic) {
+ return createApiVersionsResponse(throttleTimeMs, minMagic, Features.emptySupportedFeatures(),
+ Features.emptyFinalizedFeatures(), UNKNOWN_FINALIZED_FEATURES_EPOCH);
}
private static ApiVersionsResponse createApiVersionsResponse(
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/BeginQuorumEpochRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/BeginQuorumEpochRequest.java
index a3e492b72e233..c83e29dd3d52a 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/BeginQuorumEpochRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/BeginQuorumEpochRequest.java
@@ -20,9 +20,10 @@
import org.apache.kafka.common.message.BeginQuorumEpochRequestData;
import org.apache.kafka.common.message.BeginQuorumEpochResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
+import java.nio.ByteBuffer;
import java.util.Collections;
public class BeginQuorumEpochRequest extends AbstractRequest {
@@ -52,14 +53,9 @@ private BeginQuorumEpochRequest(BeginQuorumEpochRequestData data, short version)
this.data = data;
}
- public BeginQuorumEpochRequest(Struct struct, short version) {
- super(ApiKeys.BEGIN_QUORUM_EPOCH, version);
- this.data = new BeginQuorumEpochRequestData(struct, version);
- }
-
@Override
- protected Struct toStruct() {
- return data.toStruct(version());
+ protected BeginQuorumEpochRequestData data() {
+ return data;
}
@Override
@@ -68,6 +64,10 @@ public BeginQuorumEpochResponse getErrorResponse(int throttleTimeMs, Throwable e
.setErrorCode(Errors.forException(e).code()));
}
+ public static BeginQuorumEpochRequest parse(ByteBuffer buffer, short version) {
+ return new BeginQuorumEpochRequest(new BeginQuorumEpochRequestData(new ByteBufferAccessor(buffer), version), version);
+ }
+
public static BeginQuorumEpochRequestData singletonRequest(TopicPartition topicPartition,
int leaderEpoch,
int leaderId) {
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/BeginQuorumEpochResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/BeginQuorumEpochResponse.java
index 24eedb2bd82c6..c3e80eccaa82c 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/BeginQuorumEpochResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/BeginQuorumEpochResponse.java
@@ -20,8 +20,8 @@
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.message.BeginQuorumEpochResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.Collections;
@@ -45,18 +45,10 @@ public class BeginQuorumEpochResponse extends AbstractResponse {
public final BeginQuorumEpochResponseData data;
public BeginQuorumEpochResponse(BeginQuorumEpochResponseData data) {
+ super(ApiKeys.BEGIN_QUORUM_EPOCH);
this.data = data;
}
- public BeginQuorumEpochResponse(Struct struct, short version) {
- this.data = new BeginQuorumEpochResponseData(struct, version);
- }
-
- public BeginQuorumEpochResponse(Struct struct) {
- short latestVersion = (short) (BeginQuorumEpochResponseData.SCHEMAS.length - 1);
- this.data = new BeginQuorumEpochResponseData(struct, latestVersion);
- }
-
public static BeginQuorumEpochResponseData singletonResponse(
Errors topLevelError,
TopicPartition topicPartition,
@@ -78,11 +70,6 @@ public static BeginQuorumEpochResponseData singletonResponse(
);
}
- @Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
- }
-
@Override
public Map errorCounts() {
Map errors = new HashMap<>();
@@ -98,8 +85,18 @@ public Map errorCounts() {
return errors;
}
+ @Override
+ protected BeginQuorumEpochResponseData data() {
+ return data;
+ }
+
+ @Override
+ public int throttleTimeMs() {
+ return DEFAULT_THROTTLE_TIME;
+ }
+
public static BeginQuorumEpochResponse parse(ByteBuffer buffer, short version) {
- return new BeginQuorumEpochResponse(ApiKeys.BEGIN_QUORUM_EPOCH.responseSchema(version).read(buffer), version);
+ return new BeginQuorumEpochResponse(new BeginQuorumEpochResponseData(new ByteBufferAccessor(buffer), version));
}
-}
\ No newline at end of file
+}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ControlledShutdownRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/ControlledShutdownRequest.java
index 71238452cf782..f3c063a7a060d 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/ControlledShutdownRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/ControlledShutdownRequest.java
@@ -19,8 +19,8 @@
import org.apache.kafka.common.message.ControlledShutdownRequestData;
import org.apache.kafka.common.message.ControlledShutdownResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
@@ -47,18 +47,10 @@ public String toString() {
}
private final ControlledShutdownRequestData data;
- private final short version;
private ControlledShutdownRequest(ControlledShutdownRequestData data, short version) {
super(ApiKeys.CONTROLLED_SHUTDOWN, version);
this.data = data;
- this.version = version;
- }
-
- public ControlledShutdownRequest(Struct struct, short version) {
- super(ApiKeys.CONTROLLED_SHUTDOWN, version);
- this.data = new ControlledShutdownRequestData(struct, version);
- this.version = version;
}
@Override
@@ -69,13 +61,8 @@ public ControlledShutdownResponse getErrorResponse(int throttleTimeMs, Throwable
}
public static ControlledShutdownRequest parse(ByteBuffer buffer, short version) {
- return new ControlledShutdownRequest(
- ApiKeys.CONTROLLED_SHUTDOWN.parseRequest(version, buffer), version);
- }
-
- @Override
- protected Struct toStruct() {
- return data.toStruct(version);
+ return new ControlledShutdownRequest(new ControlledShutdownRequestData(new ByteBufferAccessor(buffer), version),
+ version);
}
public ControlledShutdownRequestData data() {
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ControlledShutdownResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/ControlledShutdownResponse.java
index 53742c21a49a6..1add8000364e0 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/ControlledShutdownResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/ControlledShutdownResponse.java
@@ -20,8 +20,8 @@
import org.apache.kafka.common.message.ControlledShutdownResponseData;
import org.apache.kafka.common.message.ControlledShutdownResponseData.RemainingPartition;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.Map;
@@ -40,13 +40,10 @@ public class ControlledShutdownResponse extends AbstractResponse {
private final ControlledShutdownResponseData data;
public ControlledShutdownResponse(ControlledShutdownResponseData data) {
+ super(ApiKeys.CONTROLLED_SHUTDOWN);
this.data = data;
}
- public ControlledShutdownResponse(Struct struct, short version) {
- this.data = new ControlledShutdownResponseData(struct, version);
- }
-
public Errors error() {
return Errors.forCode(data.errorCode());
}
@@ -56,13 +53,13 @@ public Map errorCounts() {
return errorCounts(error());
}
- public static ControlledShutdownResponse parse(ByteBuffer buffer, short version) {
- return new ControlledShutdownResponse(ApiKeys.CONTROLLED_SHUTDOWN.parseResponse(version, buffer), version);
+ @Override
+ public int throttleTimeMs() {
+ return DEFAULT_THROTTLE_TIME;
}
- @Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
+ public static ControlledShutdownResponse parse(ByteBuffer buffer, short version) {
+ return new ControlledShutdownResponse(new ControlledShutdownResponseData(new ByteBufferAccessor(buffer), version));
}
public ControlledShutdownResponseData data() {
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/CreateAclsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/CreateAclsRequest.java
index 3eb88a9b9144b..2ce651565b745 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/CreateAclsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/CreateAclsRequest.java
@@ -21,15 +21,15 @@
import org.apache.kafka.common.acl.AclBinding;
import org.apache.kafka.common.acl.AclOperation;
import org.apache.kafka.common.acl.AclPermissionType;
+import org.apache.kafka.common.errors.UnsupportedVersionException;
import org.apache.kafka.common.message.CreateAclsRequestData;
import org.apache.kafka.common.message.CreateAclsRequestData.AclCreation;
import org.apache.kafka.common.message.CreateAclsResponseData;
import org.apache.kafka.common.message.CreateAclsResponseData.AclCreationResult;
-import org.apache.kafka.common.resource.ResourcePattern;
-import org.apache.kafka.common.errors.UnsupportedVersionException;
import org.apache.kafka.common.protocol.ApiKeys;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.resource.PatternType;
+import org.apache.kafka.common.resource.ResourcePattern;
import org.apache.kafka.common.resource.ResourceType;
import java.nio.ByteBuffer;
@@ -48,7 +48,7 @@ public Builder(CreateAclsRequestData data) {
@Override
public CreateAclsRequest build(short version) {
- return new CreateAclsRequest(version, data);
+ return new CreateAclsRequest(data, version);
}
@Override
@@ -59,23 +59,19 @@ public String toString() {
private final CreateAclsRequestData data;
- CreateAclsRequest(short version, CreateAclsRequestData data) {
+ CreateAclsRequest(CreateAclsRequestData data, short version) {
super(ApiKeys.CREATE_ACLS, version);
validate(data);
this.data = data;
}
- public CreateAclsRequest(Struct struct, short version) {
- this(version, new CreateAclsRequestData(struct, version));
- }
-
public List aclCreations() {
return data.creations();
}
@Override
- protected Struct toStruct() {
- return data.toStruct(version());
+ protected CreateAclsRequestData data() {
+ return data;
}
@Override
@@ -88,7 +84,7 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable throwable
}
public static CreateAclsRequest parse(ByteBuffer buffer, short version) {
- return new CreateAclsRequest(ApiKeys.CREATE_ACLS.parseRequest(version, buffer), version);
+ return new CreateAclsRequest(new CreateAclsRequestData(new ByteBufferAccessor(buffer), version), version);
}
private void validate(CreateAclsRequestData data) {
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/CreateAclsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/CreateAclsResponse.java
index 5afce9eba217e..d0149b9a2a0d7 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/CreateAclsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/CreateAclsResponse.java
@@ -18,8 +18,8 @@
import org.apache.kafka.common.message.CreateAclsResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.List;
@@ -29,16 +29,13 @@ public class CreateAclsResponse extends AbstractResponse {
private final CreateAclsResponseData data;
public CreateAclsResponse(CreateAclsResponseData data) {
+ super(ApiKeys.CREATE_ACLS);
this.data = data;
}
- public CreateAclsResponse(Struct struct, short version) {
- this.data = new CreateAclsResponseData(struct, version);
- }
-
@Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
+ protected CreateAclsResponseData data() {
+ return data;
}
@Override
@@ -56,7 +53,7 @@ public Map errorCounts() {
}
public static CreateAclsResponse parse(ByteBuffer buffer, short version) {
- return new CreateAclsResponse(ApiKeys.CREATE_ACLS.responseSchema(version).read(buffer), version);
+ return new CreateAclsResponse(new CreateAclsResponseData(new ByteBufferAccessor(buffer), version));
}
@Override
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/CreateDelegationTokenRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/CreateDelegationTokenRequest.java
index 61b404a7e0623..e1d4cfaab80bd 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/CreateDelegationTokenRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/CreateDelegationTokenRequest.java
@@ -18,8 +18,8 @@
import org.apache.kafka.common.message.CreateDelegationTokenRequestData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import org.apache.kafka.common.security.auth.KafkaPrincipal;
import java.nio.ByteBuffer;
@@ -33,18 +33,9 @@ private CreateDelegationTokenRequest(CreateDelegationTokenRequestData data, shor
this.data = data;
}
- public CreateDelegationTokenRequest(Struct struct, short version) {
- super(ApiKeys.CREATE_DELEGATION_TOKEN, version);
- this.data = new CreateDelegationTokenRequestData(struct, version);
- }
-
public static CreateDelegationTokenRequest parse(ByteBuffer buffer, short version) {
- return new CreateDelegationTokenRequest(ApiKeys.CREATE_DELEGATION_TOKEN.parseRequest(version, buffer), version);
- }
-
- @Override
- protected Struct toStruct() {
- return data.toStruct(version());
+ return new CreateDelegationTokenRequest(new CreateDelegationTokenRequestData(new ByteBufferAccessor(buffer), version),
+ version);
}
public CreateDelegationTokenRequestData data() {
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/CreateDelegationTokenResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/CreateDelegationTokenResponse.java
index 9796d68a82b99..9d39b6b641345 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/CreateDelegationTokenResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/CreateDelegationTokenResponse.java
@@ -18,8 +18,8 @@
import org.apache.kafka.common.message.CreateDelegationTokenResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import org.apache.kafka.common.security.auth.KafkaPrincipal;
import java.nio.ByteBuffer;
@@ -30,15 +30,13 @@ public class CreateDelegationTokenResponse extends AbstractResponse {
private final CreateDelegationTokenResponseData data;
public CreateDelegationTokenResponse(CreateDelegationTokenResponseData data) {
+ super(ApiKeys.CREATE_DELEGATION_TOKEN);
this.data = data;
}
- public CreateDelegationTokenResponse(Struct struct, short version) {
- this.data = new CreateDelegationTokenResponseData(struct, version);
- }
-
public static CreateDelegationTokenResponse parse(ByteBuffer buffer, short version) {
- return new CreateDelegationTokenResponse(ApiKeys.CREATE_DELEGATION_TOKEN.responseSchema(version).read(buffer), version);
+ return new CreateDelegationTokenResponse(
+ new CreateDelegationTokenResponseData(new ByteBufferAccessor(buffer), version));
}
public static CreateDelegationTokenResponse prepareResponse(int throttleTimeMs,
@@ -75,11 +73,6 @@ public Map errorCounts() {
return errorCounts(error());
}
- @Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
- }
-
@Override
public int throttleTimeMs() {
return data.throttleTimeMs();
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/CreatePartitionsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/CreatePartitionsRequest.java
index 77d4175c339de..8b0b9f81cc92c 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/CreatePartitionsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/CreatePartitionsRequest.java
@@ -22,7 +22,7 @@
import org.apache.kafka.common.message.CreatePartitionsResponseData;
import org.apache.kafka.common.message.CreatePartitionsResponseData.CreatePartitionsTopicResult;
import org.apache.kafka.common.protocol.ApiKeys;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import java.nio.ByteBuffer;
@@ -55,15 +55,6 @@ public String toString() {
this.data = data;
}
- public CreatePartitionsRequest(Struct struct, short apiVersion) {
- this(new CreatePartitionsRequestData(struct, apiVersion), apiVersion);
- }
-
- @Override
- protected Struct toStruct() {
- return data.toStruct(version());
- }
-
public CreatePartitionsRequestData data() {
return data;
}
@@ -85,6 +76,6 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
}
public static CreatePartitionsRequest parse(ByteBuffer buffer, short version) {
- return new CreatePartitionsRequest(ApiKeys.CREATE_PARTITIONS.parseRequest(version, buffer), version);
+ return new CreatePartitionsRequest(new CreatePartitionsRequestData(new ByteBufferAccessor(buffer), version), version);
}
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/CreatePartitionsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/CreatePartitionsResponse.java
index d658b0e4c230d..e0af04b07ae6a 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/CreatePartitionsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/CreatePartitionsResponse.java
@@ -19,35 +19,26 @@
import org.apache.kafka.common.message.CreatePartitionsResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.HashMap;
import java.util.Map;
-
public class CreatePartitionsResponse extends AbstractResponse {
private final CreatePartitionsResponseData data;
public CreatePartitionsResponse(CreatePartitionsResponseData data) {
+ super(ApiKeys.CREATE_PARTITIONS);
this.data = data;
}
- public CreatePartitionsResponse(Struct struct, short version) {
- this.data = new CreatePartitionsResponseData(struct, version);
- }
-
public CreatePartitionsResponseData data() {
return data;
}
- @Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
- }
-
@Override
public Map errorCounts() {
Map counts = new HashMap<>();
@@ -58,7 +49,7 @@ public Map errorCounts() {
}
public static CreatePartitionsResponse parse(ByteBuffer buffer, short version) {
- return new CreatePartitionsResponse(ApiKeys.CREATE_PARTITIONS.parseResponse(version, buffer), version);
+ return new CreatePartitionsResponse(new CreatePartitionsResponseData(new ByteBufferAccessor(buffer), version));
}
@Override
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/CreateTopicsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/CreateTopicsRequest.java
index dd26e5642632e..7ec2c9adde3c4 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/CreateTopicsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/CreateTopicsRequest.java
@@ -25,7 +25,7 @@
import org.apache.kafka.common.message.CreateTopicsResponseData;
import org.apache.kafka.common.message.CreateTopicsResponseData.CreatableTopicResult;
import org.apache.kafka.common.protocol.ApiKeys;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
public class CreateTopicsRequest extends AbstractRequest {
public static class Builder extends AbstractRequest.Builder {
@@ -72,16 +72,11 @@ public String toString() {
public static final int NO_NUM_PARTITIONS = -1;
public static final short NO_REPLICATION_FACTOR = -1;
- private CreateTopicsRequest(CreateTopicsRequestData data, short version) {
+ public CreateTopicsRequest(CreateTopicsRequestData data, short version) {
super(ApiKeys.CREATE_TOPICS, version);
this.data = data;
}
- public CreateTopicsRequest(Struct struct, short version) {
- super(ApiKeys.CREATE_TOPICS, version);
- this.data = new CreateTopicsRequestData(struct, version);
- }
-
public CreateTopicsRequestData data() {
return data;
}
@@ -103,14 +98,6 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
}
public static CreateTopicsRequest parse(ByteBuffer buffer, short version) {
- return new CreateTopicsRequest(ApiKeys.CREATE_TOPICS.parseRequest(version, buffer), version);
- }
-
- /**
- * Visible for testing.
- */
- @Override
- public Struct toStruct() {
- return data.toStruct(version());
+ return new CreateTopicsRequest(new CreateTopicsRequestData(new ByteBufferAccessor(buffer), version), version);
}
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/CreateTopicsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/CreateTopicsResponse.java
index 7b64684d4dd4a..c15da20dfc79a 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/CreateTopicsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/CreateTopicsResponse.java
@@ -19,8 +19,8 @@
import org.apache.kafka.common.message.CreateTopicsResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.HashMap;
@@ -46,22 +46,14 @@ public class CreateTopicsResponse extends AbstractResponse {
private final CreateTopicsResponseData data;
public CreateTopicsResponse(CreateTopicsResponseData data) {
+ super(ApiKeys.CREATE_TOPICS);
this.data = data;
}
- public CreateTopicsResponse(Struct struct, short version) {
- this.data = new CreateTopicsResponseData(struct, version);
- }
-
public CreateTopicsResponseData data() {
return data;
}
- @Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
- }
-
@Override
public int throttleTimeMs() {
return data.throttleTimeMs();
@@ -77,8 +69,7 @@ public Map errorCounts() {
}
public static CreateTopicsResponse parse(ByteBuffer buffer, short version) {
- return new CreateTopicsResponse(
- ApiKeys.CREATE_TOPICS.responseSchema(version).read(buffer), version);
+ return new CreateTopicsResponse(new CreateTopicsResponseData(new ByteBufferAccessor(buffer), version));
}
@Override
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DeleteAclsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/DeleteAclsRequest.java
index b020b2b0ae3d3..6face08f5675d 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/DeleteAclsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/DeleteAclsRequest.java
@@ -28,7 +28,7 @@
import org.apache.kafka.common.message.DeleteAclsResponseData;
import org.apache.kafka.common.message.DeleteAclsResponseData.DeleteAclsFilterResult;
import org.apache.kafka.common.protocol.ApiKeys;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.resource.PatternType;
import org.apache.kafka.common.resource.ResourcePatternFilter;
import org.apache.kafka.common.resource.ResourceType;
@@ -49,7 +49,7 @@ public Builder(DeleteAclsRequestData data) {
@Override
public DeleteAclsRequest build(short version) {
- return new DeleteAclsRequest(version, data);
+ return new DeleteAclsRequest(data, version);
}
@Override
@@ -61,7 +61,7 @@ public String toString() {
private final DeleteAclsRequestData data;
- private DeleteAclsRequest(short version, DeleteAclsRequestData data) {
+ private DeleteAclsRequest(DeleteAclsRequestData data, short version) {
super(ApiKeys.DELETE_ACLS, version);
this.data = data;
normalizeAndValidate();
@@ -95,18 +95,13 @@ else if (patternType != PatternType.LITERAL)
}
}
- public DeleteAclsRequest(Struct struct, short version) {
- super(ApiKeys.DELETE_ACLS, version);
- this.data = new DeleteAclsRequestData(struct, version);
- }
-
public List filters() {
return data.filters().stream().map(DeleteAclsRequest::aclBindingFilter).collect(Collectors.toList());
}
@Override
- protected Struct toStruct() {
- return data.toStruct(version());
+ protected DeleteAclsRequestData data() {
+ return data;
}
@Override
@@ -118,11 +113,11 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable throwable
.setErrorMessage(apiError.message()));
return new DeleteAclsResponse(new DeleteAclsResponseData()
.setThrottleTimeMs(throttleTimeMs)
- .setFilterResults(filterResults));
+ .setFilterResults(filterResults), version());
}
public static DeleteAclsRequest parse(ByteBuffer buffer, short version) {
- return new DeleteAclsRequest(DELETE_ACLS.parseRequest(version, buffer), version);
+ return new DeleteAclsRequest(new DeleteAclsRequestData(new ByteBufferAccessor(buffer), version), version);
}
public static DeleteAclsFilter deleteAclsFilter(AclBindingFilter filter) {
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DeleteAclsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DeleteAclsResponse.java
index 33d72a1df77fb..3ff8a9834f888 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/DeleteAclsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/DeleteAclsResponse.java
@@ -18,17 +18,17 @@
import org.apache.kafka.common.acl.AccessControlEntry;
import org.apache.kafka.common.acl.AclBinding;
+import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
+import org.apache.kafka.common.resource.ResourcePattern;
import org.apache.kafka.common.acl.AclOperation;
import org.apache.kafka.common.acl.AclPermissionType;
import org.apache.kafka.common.errors.UnsupportedVersionException;
import org.apache.kafka.common.message.DeleteAclsResponseData;
import org.apache.kafka.common.message.DeleteAclsResponseData.DeleteAclsFilterResult;
import org.apache.kafka.common.message.DeleteAclsResponseData.DeleteAclsMatchingAcl;
-import org.apache.kafka.common.protocol.ApiKeys;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import org.apache.kafka.common.resource.PatternType;
-import org.apache.kafka.common.resource.ResourcePattern;
import org.apache.kafka.common.resource.ResourceType;
import org.apache.kafka.server.authorizer.AclDeleteResult;
import org.slf4j.Logger;
@@ -44,18 +44,15 @@ public class DeleteAclsResponse extends AbstractResponse {
private final DeleteAclsResponseData data;
- public DeleteAclsResponse(DeleteAclsResponseData data) {
+ public DeleteAclsResponse(DeleteAclsResponseData data, short version) {
+ super(ApiKeys.DELETE_ACLS);
this.data = data;
- }
-
- public DeleteAclsResponse(Struct struct, short version) {
- data = new DeleteAclsResponseData(struct, version);
+ validate(version);
}
@Override
- protected Struct toStruct(short version) {
- validate(version);
- return data.toStruct(version);
+ protected DeleteAclsResponseData data() {
+ return data;
}
@Override
@@ -73,7 +70,7 @@ public Map errorCounts() {
}
public static DeleteAclsResponse parse(ByteBuffer buffer, short version) {
- return new DeleteAclsResponse(ApiKeys.DELETE_ACLS.parseResponse(version, buffer), version);
+ return new DeleteAclsResponse(new DeleteAclsResponseData(new ByteBufferAccessor(buffer), version), version);
}
public String toString() {
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DeleteGroupsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/DeleteGroupsRequest.java
index 2ad947c5904be..d09a4d4c02d77 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/DeleteGroupsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/DeleteGroupsRequest.java
@@ -21,8 +21,8 @@
import org.apache.kafka.common.message.DeleteGroupsResponseData.DeletableGroupResult;
import org.apache.kafka.common.message.DeleteGroupsResponseData.DeletableGroupResultCollection;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
@@ -53,11 +53,6 @@ public DeleteGroupsRequest(DeleteGroupsRequestData data, short version) {
this.data = data;
}
- public DeleteGroupsRequest(Struct struct, short version) {
- super(ApiKeys.DELETE_GROUPS, version);
- this.data = new DeleteGroupsRequestData(struct, version);
- }
-
@Override
public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
Errors error = Errors.forException(e);
@@ -76,11 +71,11 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
}
public static DeleteGroupsRequest parse(ByteBuffer buffer, short version) {
- return new DeleteGroupsRequest(ApiKeys.DELETE_GROUPS.parseRequest(version, buffer), version);
+ return new DeleteGroupsRequest(new DeleteGroupsRequestData(new ByteBufferAccessor(buffer), version), version);
}
@Override
- protected Struct toStruct() {
- return data.toStruct(version());
+ protected DeleteGroupsRequestData data() {
+ return data;
}
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DeleteGroupsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DeleteGroupsResponse.java
index ad171c8f9d03b..a8e8d482de1d1 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/DeleteGroupsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/DeleteGroupsResponse.java
@@ -19,8 +19,8 @@
import org.apache.kafka.common.message.DeleteGroupsResponseData;
import org.apache.kafka.common.message.DeleteGroupsResponseData.DeletableGroupResult;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.HashMap;
@@ -42,21 +42,13 @@ public class DeleteGroupsResponse extends AbstractResponse {
public final DeleteGroupsResponseData data;
public DeleteGroupsResponse(DeleteGroupsResponseData data) {
+ super(ApiKeys.DELETE_GROUPS);
this.data = data;
}
- public DeleteGroupsResponse(Struct struct) {
- short latestVersion = (short) (DeleteGroupsResponseData.SCHEMAS.length - 1);
- this.data = new DeleteGroupsResponseData(struct, latestVersion);
- }
-
- public DeleteGroupsResponse(Struct struct, short version) {
- this.data = new DeleteGroupsResponseData(struct, version);
- }
-
@Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
+ protected DeleteGroupsResponseData data() {
+ return data;
}
public Map errors() {
@@ -85,7 +77,7 @@ public Map errorCounts() {
}
public static DeleteGroupsResponse parse(ByteBuffer buffer, short version) {
- return new DeleteGroupsResponse(ApiKeys.DELETE_GROUPS.parseResponse(version, buffer), version);
+ return new DeleteGroupsResponse(new DeleteGroupsResponseData(new ByteBufferAccessor(buffer), version));
}
@Override
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DeleteRecordsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/DeleteRecordsRequest.java
index 1939923c554e4..a4f62e18f675a 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/DeleteRecordsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/DeleteRecordsRequest.java
@@ -22,8 +22,8 @@
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.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
@@ -57,16 +57,7 @@ private DeleteRecordsRequest(DeleteRecordsRequestData data, short version) {
this.data = data;
}
- public DeleteRecordsRequest(Struct struct, short version) {
- super(ApiKeys.DELETE_RECORDS, version);
- this.data = new DeleteRecordsRequestData(struct, version);
- }
-
@Override
- protected Struct toStruct() {
- return data.toStruct(version());
- }
-
public DeleteRecordsRequestData data() {
return data;
}
@@ -89,6 +80,6 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
}
public static DeleteRecordsRequest parse(ByteBuffer buffer, short version) {
- return new DeleteRecordsRequest(ApiKeys.DELETE_RECORDS.parseRequest(version, buffer), version);
+ return new DeleteRecordsRequest(new DeleteRecordsRequestData(new ByteBufferAccessor(buffer), version), version);
}
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DeleteRecordsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DeleteRecordsResponse.java
index 968a35bd32c59..ef34102f8a8fd 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/DeleteRecordsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/DeleteRecordsResponse.java
@@ -19,8 +19,8 @@
import org.apache.kafka.common.message.DeleteRecordsResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.HashMap;
@@ -42,18 +42,10 @@ public class DeleteRecordsResponse extends AbstractResponse {
*/
public DeleteRecordsResponse(DeleteRecordsResponseData data) {
+ super(ApiKeys.DELETE_RECORDS);
this.data = data;
}
- public DeleteRecordsResponse(Struct struct, short version) {
- this.data = new DeleteRecordsResponseData(struct, version);
- }
-
- @Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
- }
-
public DeleteRecordsResponseData data() {
return data;
}
@@ -75,7 +67,7 @@ public Map errorCounts() {
}
public static DeleteRecordsResponse parse(ByteBuffer buffer, short version) {
- return new DeleteRecordsResponse(ApiKeys.DELETE_RECORDS.parseResponse(version, buffer), version);
+ return new DeleteRecordsResponse(new DeleteRecordsResponseData(new ByteBufferAccessor(buffer), version));
}
@Override
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DeleteTopicsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/DeleteTopicsRequest.java
index 440acfd2b838a..5ab64184aeab0 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/DeleteTopicsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/DeleteTopicsRequest.java
@@ -20,15 +20,12 @@
import org.apache.kafka.common.message.DeleteTopicsResponseData;
import org.apache.kafka.common.message.DeleteTopicsResponseData.DeletableTopicResult;
import org.apache.kafka.common.protocol.ApiKeys;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import java.nio.ByteBuffer;
public class DeleteTopicsRequest extends AbstractRequest {
- private DeleteTopicsRequestData data;
- private final short version;
-
public static class Builder extends AbstractRequest.Builder {
private DeleteTopicsRequestData data;
@@ -48,21 +45,11 @@ public String toString() {
}
}
+ private DeleteTopicsRequestData data;
+
private DeleteTopicsRequest(DeleteTopicsRequestData data, short version) {
super(ApiKeys.DELETE_TOPICS, version);
this.data = data;
- this.version = version;
- }
-
- public DeleteTopicsRequest(Struct struct, short version) {
- super(ApiKeys.DELETE_TOPICS, version);
- this.data = new DeleteTopicsRequestData(struct, version);
- this.version = version;
- }
-
- @Override
- protected Struct toStruct() {
- return data.toStruct(version);
}
public DeleteTopicsRequestData data() {
@@ -72,7 +59,7 @@ public DeleteTopicsRequestData data() {
@Override
public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
DeleteTopicsResponseData response = new DeleteTopicsResponseData();
- if (version >= 1) {
+ if (version() >= 1) {
response.setThrottleTimeMs(throttleTimeMs);
}
ApiError apiError = ApiError.fromThrowable(e);
@@ -85,7 +72,7 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
}
public static DeleteTopicsRequest parse(ByteBuffer buffer, short version) {
- return new DeleteTopicsRequest(ApiKeys.DELETE_TOPICS.parseRequest(version, buffer), version);
+ return new DeleteTopicsRequest(new DeleteTopicsRequestData(new ByteBufferAccessor(buffer), version), version);
}
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DeleteTopicsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DeleteTopicsResponse.java
index bad23f6a2bd0b..3584ddfbfc497 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/DeleteTopicsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/DeleteTopicsResponse.java
@@ -18,8 +18,8 @@
import org.apache.kafka.common.message.DeleteTopicsResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import java.nio.ByteBuffer;
import java.util.HashMap;
@@ -41,18 +41,10 @@ public class DeleteTopicsResponse extends AbstractResponse {
private DeleteTopicsResponseData data;
public DeleteTopicsResponse(DeleteTopicsResponseData data) {
+ super(ApiKeys.DELETE_TOPICS);
this.data = data;
}
- public DeleteTopicsResponse(Struct struct, short version) {
- this.data = new DeleteTopicsResponseData(struct, version);
- }
-
- @Override
- protected Struct toStruct(short version) {
- return data.toStruct(version);
- }
-
@Override
public int throttleTimeMs() {
return data.throttleTimeMs();
@@ -72,7 +64,7 @@ public Map errorCounts() {
}
public static DeleteTopicsResponse parse(ByteBuffer buffer, short version) {
- return new DeleteTopicsResponse(ApiKeys.DELETE_TOPICS.parseResponse(version, buffer), version);
+ return new DeleteTopicsResponse(new DeleteTopicsResponseData(new ByteBufferAccessor(buffer), version));
}
@Override
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java
index 98059cfddbb8d..e6f9c38fd598b 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java
@@ -24,7 +24,7 @@
import org.apache.kafka.common.message.DescribeAclsRequestData;
import org.apache.kafka.common.message.DescribeAclsResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.resource.PatternType;
import org.apache.kafka.common.resource.ResourcePatternFilter;
import org.apache.kafka.common.resource.ResourceType;
@@ -89,20 +89,10 @@ else if (patternType != PatternType.LITERAL)
}
}
- public DescribeAclsRequest(Struct struct, short version) {
- super(ApiKeys.DESCRIBE_ACLS, version);
- this.data = new DescribeAclsRequestData(struct, version);
- }
-
public DescribeAclsRequestData data() {
return data;
}
- @Override
- protected Struct toStruct() {
- return data.toStruct(version());
- }
-
@Override
public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable throwable) {
ApiError error = ApiError.fromThrowable(throwable);
@@ -110,11 +100,11 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable throwable
.setThrottleTimeMs(throttleTimeMs)
.setErrorCode(error.error().code())
.setErrorMessage(error.message());
- return new DescribeAclsResponse(response);
+ return new DescribeAclsResponse(response, version());
}
public static DescribeAclsRequest parse(ByteBuffer buffer, short version) {
- return new DescribeAclsRequest(ApiKeys.DESCRIBE_ACLS.parseRequest(version, buffer), version);
+ return new DescribeAclsRequest(new DescribeAclsRequestData(new ByteBufferAccessor(buffer), version), version);
}
public AclBindingFilter filter() {
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java
index 3a41cb3de85b1..4308c9ea45bf7 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java
@@ -17,6 +17,15 @@
package org.apache.kafka.common.requests;
+import org.apache.kafka.common.acl.AccessControlEntry;
+import org.apache.kafka.common.acl.AclBinding;
+import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
+import org.apache.kafka.common.resource.PatternType;
+import org.apache.kafka.common.resource.ResourcePattern;
+import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.protocol.Errors;
+
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Collection;
@@ -24,40 +33,37 @@
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
+import java.util.Optional;
import java.util.stream.Collectors;
import java.util.stream.Stream;
-import org.apache.kafka.common.acl.AccessControlEntry;
-import org.apache.kafka.common.acl.AclBinding;
import org.apache.kafka.common.acl.AclOperation;
import org.apache.kafka.common.acl.AclPermissionType;
-import org.apache.kafka.common.errors.UnsupportedVersionException;
import org.apache.kafka.common.message.DescribeAclsResponseData;
import org.apache.kafka.common.message.DescribeAclsResponseData.AclDescription;
import org.apache.kafka.common.message.DescribeAclsResponseData.DescribeAclsResource;
-import org.apache.kafka.common.protocol.ApiKeys;
-import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
-import org.apache.kafka.common.resource.PatternType;
-import org.apache.kafka.common.resource.ResourcePattern;
import org.apache.kafka.common.resource.ResourceType;
public class DescribeAclsResponse extends AbstractResponse {
private final DescribeAclsResponseData data;
- public DescribeAclsResponse(DescribeAclsResponseData data) {
+ public DescribeAclsResponse(DescribeAclsResponseData data, short version) {
+ super(ApiKeys.DESCRIBE_ACLS);
this.data = data;
+ validate(Optional.of(version));
}
- public DescribeAclsResponse(Struct struct, short version) {
- this.data = new DescribeAclsResponseData(struct, version);
+ // Skips version validation, visible for testing
+ DescribeAclsResponse(DescribeAclsResponseData data) {
+ super(ApiKeys.DESCRIBE_ACLS);
+ this.data = data;
+ validate(Optional.empty());
}
@Override
- protected Struct toStruct(short version) {
- validate(version);
- return data.toStruct(version);
+ protected DescribeAclsResponseData data() {
+ return data;
}
@Override
@@ -79,7 +85,7 @@ public List acls() {
}
public static DescribeAclsResponse parse(ByteBuffer buffer, short version) {
- return new DescribeAclsResponse(ApiKeys.DESCRIBE_ACLS.responseSchema(version).read(buffer), version);
+ return new DescribeAclsResponse(new DescribeAclsResponseData(new ByteBufferAccessor(buffer), version), version);
}
@Override
@@ -87,8 +93,8 @@ public boolean shouldClientThrottle(short version) {
return version >= 1;
}
- private void validate(short version) {
- if (version == 0) {
+ private void validate(Optional version) {
+ if (version.isPresent() && version.get() == 0) {
final boolean unsupported = acls().stream()
.anyMatch(acl -> acl.patternType() != PatternType.LITERAL.code());
if (unsupported) {
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeClientQuotasRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeClientQuotasRequest.java
index 593cc6a9aace0..1f167b0ee3c2b 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeClientQuotasRequest.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeClientQuotasRequest.java
@@ -18,11 +18,13 @@
import org.apache.kafka.common.message.DescribeClientQuotasRequestData;
import org.apache.kafka.common.message.DescribeClientQuotasRequestData.ComponentData;
+import org.apache.kafka.common.message.DescribeClientQuotasResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
-import org.apache.kafka.common.protocol.types.Struct;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.quota.ClientQuotaFilter;
import org.apache.kafka.common.quota.ClientQuotaFilterComponent;
+import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.List;
@@ -77,11 +79,6 @@ public DescribeClientQuotasRequest(DescribeClientQuotasRequestData data, short v
this.data = data;
}
- public DescribeClientQuotasRequest(Struct struct, short version) {
- super(ApiKeys.DESCRIBE_CLIENT_QUOTAS, version);
- this.data = new DescribeClientQuotasRequestData(struct, version);
- }
-
public ClientQuotaFilter filter() {
List components = new ArrayList<>(data.components().size());
for (ComponentData componentData : data.components()) {
@@ -109,12 +106,23 @@ public ClientQuotaFilter filter() {
}
@Override
- public DescribeClientQuotasResponse getErrorResponse(int throttleTimeMs, Throwable e) {
- return new DescribeClientQuotasResponse(throttleTimeMs, e);
+ protected DescribeClientQuotasRequestData data() {
+ return data;
}
@Override
- protected Struct toStruct() {
- return data.toStruct(version());
+ public DescribeClientQuotasResponse getErrorResponse(int throttleTimeMs, Throwable e) {
+ ApiError error = ApiError.fromThrowable(e);
+ return new DescribeClientQuotasResponse(new DescribeClientQuotasResponseData()
+ .setThrottleTimeMs(throttleTimeMs)
+ .setErrorCode(error.error().code())
+ .setErrorMessage(error.message())
+ .setEntries(null));
+ }
+
+ public static DescribeClientQuotasRequest parse(ByteBuffer buffer, short version) {
+ return new DescribeClientQuotasRequest(new DescribeClientQuotasRequestData(new ByteBufferAccessor(buffer), version),
+ version);
}
+
}
diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeClientQuotasResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeClientQuotasResponse.java
index bda3673c27f2c..94fca64221d0f 100644
--- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeClientQuotasResponse.java
+++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeClientQuotasResponse.java
@@ -22,8 +22,8 @@
import org.apache.kafka.common.message.DescribeClientQuotasResponseData.EntryData;
import org.apache.kafka.common.message.DescribeClientQuotasResponseData.ValueData;
import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.ByteBufferAccessor;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.common.protocol.types.Struct;
import org.apache.kafka.common.quota.ClientQuotaEntity;
import java.nio.ByteBuffer;
@@ -36,48 +36,9 @@ public class DescribeClientQuotasResponse extends AbstractResponse {
private final DescribeClientQuotasResponseData data;
- public DescribeClientQuotasResponse(Map> entities, int throttleTimeMs) {
- List entries = new ArrayList<>(entities.size());
- for (Map.Entry> entry : entities.entrySet()) {
- ClientQuotaEntity quotaEntity = entry.getKey();
- List entityData = new ArrayList<>(quotaEntity.entries().size());
- for (Map.Entry entityEntry : quotaEntity.entries().entrySet()) {
- entityData.add(new EntityData()
- .setEntityType(entityEntry.getKey())
- .setEntityName(entityEntry.getValue()));
- }
-
- Map quotaValues = entry.getValue();
- List valueData = new ArrayList<>(quotaValues.size());
- for (Map.Entry valuesEntry : entry.getValue().entrySet()) {
- valueData.add(new ValueData()
- .setKey(valuesEntry.getKey())
- .setValue(valuesEntry.getValue()));
- }
-
- entries.add(new EntryData()
- .setEntity(entityData)
- .setValues(valueData));
- }
-
- this.data = new DescribeClientQuotasResponseData()
- .setThrottleTimeMs(throttleTimeMs)
- .setErrorCode((short) 0)
- .setErrorMessage(null)
- .setEntries(entries);
- }
-
- public DescribeClientQuotasResponse(int throttleTimeMs, Throwable e) {
- ApiError apiError = ApiError.fromThrowable(e);
- this.data = new DescribeClientQuotasResponseData()
- .setThrottleTimeMs(throttleTimeMs)
- .setErrorCode(apiError.error().code())
- .setErrorMessage(apiError.message())
- .setEntries(null);
- }
-
- public DescribeClientQuotasResponse(Struct struct, short version) {
- this.data = new DescribeClientQuotasResponseData(struct, version);
+ public DescribeClientQuotasResponse(DescribeClientQuotasResponseData data) {
+ super(ApiKeys.DESCRIBE_CLIENT_QUOTAS);
+ this.data = data;
}
public void complete(KafkaFutureImpl