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 cd99f472ebb0a..7e4425d3e79c1 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 @@ -266,8 +266,20 @@ public ApiKeys apiKey() { return apiKey; } + /** + * Get the throttle time in milliseconds. If the response schema does not + * support this field, then 0 will be returned. + */ public abstract int throttleTimeMs(); + /** + * Set the throttle time in the response if the schema supports it. Otherwise, + * this is a no-op. + * + * @param throttleTimeMs The throttle time in milliseconds + */ + public abstract void maybeSetThrottleTimeMs(int throttleTimeMs); + public String toString() { return data().toString(); } 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 ce9a6cf7d6063..d90afd04ddcde 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 @@ -56,6 +56,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public AddOffsetsToTxnResponseData data() { return data; 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 57b2a5a5d7c08..8038f4b8fc66d 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 @@ -94,6 +94,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public Map errors() { if (cachedErrorsMap != null) { return cachedErrorsMap; diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AllocateProducerIdsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AllocateProducerIdsResponse.java index 41db29158e5e7..2511e2b2db320 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/AllocateProducerIdsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/AllocateProducerIdsResponse.java @@ -56,6 +56,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public Errors error() { return Errors.forCode(data.errorCode()); } 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 fcacc5d95ef07..fc56db7e73658 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 @@ -67,6 +67,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { Map counts = new HashMap<>(); 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 1115f06ee80a9..1668c2446bc77 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 @@ -55,6 +55,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public AlterConfigsResponseData data() { return data; 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 ab166b812718c..7d6c340fd147b 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 @@ -55,6 +55,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { Map counts = new HashMap<>(); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterPartitionResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterPartitionResponse.java index d2ace4112f4c1..9ee92f7b809cd 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/AlterPartitionResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterPartitionResponse.java @@ -55,6 +55,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public static AlterPartitionResponse parse(ByteBuffer buffer, short version) { return new AlterPartitionResponse(new AlterPartitionResponseData(new ByteBufferAccessor(buffer), 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 afa658d1e150f..0c38a83ee3d7e 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 @@ -53,6 +53,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { Map errorCounts = new HashMap<>(); 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 97c0b7d17b204..86c9b006a2ce0 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 @@ -48,6 +48,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { return errorCounts(data.results().stream().map(r -> Errors.forCode(r.errorCode()))); 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 28c9d613bbb26..a903e50b15d9e 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 @@ -74,6 +74,11 @@ public int throttleTimeMs() { return this.data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public boolean shouldClientThrottle(short version) { return version >= 2; 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 c8c0328c93a15..5ae975acd8a05 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 @@ -95,6 +95,11 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + // Not supported by the response schema + } + public static BeginQuorumEpochResponse parse(ByteBuffer buffer, short version) { return new BeginQuorumEpochResponse(new BeginQuorumEpochResponseData(new ByteBufferAccessor(buffer), version)); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/BrokerHeartbeatResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/BrokerHeartbeatResponse.java index e7d01e53c67ba..4c8b3aafc4dd2 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/BrokerHeartbeatResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/BrokerHeartbeatResponse.java @@ -44,6 +44,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { Map errorCounts = new HashMap<>(); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/BrokerRegistrationResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/BrokerRegistrationResponse.java index 8296d7a4c35cb..8b6121c376339 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/BrokerRegistrationResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/BrokerRegistrationResponse.java @@ -44,6 +44,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { Map errorCounts = new HashMap<>(); 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 73b6a50268379..bc5aa0ba35a33 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 @@ -58,6 +58,11 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + // Not supported by the response schema + } + public static ControlledShutdownResponse parse(ByteBuffer buffer, short version) { return new ControlledShutdownResponse(new ControlledShutdownResponseData(new ByteBufferAccessor(buffer), version)); } 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 8bc6643f9de01..cef7b73ac27e9 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 @@ -43,6 +43,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public List results() { return data.results(); } 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 22c2e1259019b..0a9f9a8991bdc 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 @@ -86,6 +86,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public Errors error() { return Errors.forCode(data.errorCode()); } 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 e59ac981f112a..2dcd2b200cadd 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 @@ -62,4 +62,9 @@ public boolean shouldClientThrottle(short version) { public int throttleTimeMs() { return data.throttleTimeMs(); } + + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } } 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 dd0627742587c..da011e224ed08 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 @@ -60,6 +60,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { HashMap counts = new HashMap<>(); 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 7482953a00d64..d0b596ed91ab1 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 @@ -60,6 +60,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public List filterResults() { return data.filterResults(); } 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 4cbffda422138..3bbb08d59fb9b 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 @@ -85,6 +85,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public boolean shouldClientThrottle(short version) { return version >= 1; 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 b090543faddfb..5084681f5373f 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 @@ -56,6 +56,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { Map errorCounts = new HashMap<>(); 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 2090c4fd2e2ee..65a54481ba07e 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 @@ -50,6 +50,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public DeleteTopicsResponseData data() { return data; 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 c4190e65640ea..c602dab9951fc 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 @@ -71,6 +71,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public ApiError error() { return new ApiError(Errors.forCode(data.errorCode()), data.errorMessage()); } 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 147414333ea08..3a052c9fe8eba 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 @@ -70,6 +70,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public DescribeClientQuotasResponseData data() { return data; diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeClusterResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeClusterResponse.java index 60d931196a659..fb48476e2467d 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeClusterResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeClusterResponse.java @@ -52,6 +52,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public DescribeClusterResponseData data() { return data; diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeConfigsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeConfigsResponse.java index aa7a713e8ab4e..6fd5320d8820e 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeConfigsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeConfigsResponse.java @@ -255,6 +255,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { Map errorCounts = new HashMap<>(); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeDelegationTokenResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeDelegationTokenResponse.java index 4fd1d99652661..a922f056a89aa 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeDelegationTokenResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeDelegationTokenResponse.java @@ -96,6 +96,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public Errors error() { return Errors.forCode(data.errorCode()); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeGroupsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeGroupsResponse.java index 360caf01e468e..119bedfdb19dc 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeGroupsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeGroupsResponse.java @@ -115,6 +115,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public static final String UNKNOWN_STATE = ""; public static final String UNKNOWN_PROTOCOL_TYPE = ""; public static final String UNKNOWN_PROTOCOL = ""; diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeLogDirsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeLogDirsResponse.java index fe8aebbc4f6b8..cbf3054217363 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeLogDirsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeLogDirsResponse.java @@ -50,6 +50,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { Map errorCounts = new HashMap<>(); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeProducersResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeProducersResponse.java index 74e9437472bc9..065a101bed6e8 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeProducersResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeProducersResponse.java @@ -66,4 +66,9 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java index 9f58e52970c47..39e050c94052b 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java @@ -70,6 +70,11 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + // Not supported by the response schema + } + public static DescribeQuorumResponseData singletonErrorResponse( TopicPartition topicPartition, Errors error diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeTransactionsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeTransactionsResponse.java index cf151b35bba35..5eef63b0ce4bb 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeTransactionsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeTransactionsResponse.java @@ -64,4 +64,10 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + } + diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeUserScramCredentialsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeUserScramCredentialsResponse.java index 001cefae41a6b..58ba4212949c6 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeUserScramCredentialsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeUserScramCredentialsResponse.java @@ -48,6 +48,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { return errorCounts(data.results().stream().map(r -> Errors.forCode(r.errorCode()))); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ElectLeadersResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/ElectLeadersResponse.java index 88d4d19fc021e..2e82cd4c9a5ef 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/ElectLeadersResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/ElectLeadersResponse.java @@ -62,6 +62,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { HashMap counts = new HashMap<>(); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/EndQuorumEpochResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/EndQuorumEpochResponse.java index ac2c0c5c9d5c3..b3a236adc69cb 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/EndQuorumEpochResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/EndQuorumEpochResponse.java @@ -73,6 +73,11 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + // Not supported by the response schema + } + public static EndQuorumEpochResponseData singletonResponse( Errors topLevelError, TopicPartition topicPartition, diff --git a/clients/src/main/java/org/apache/kafka/common/requests/EndTxnResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/EndTxnResponse.java index 029e7d0ce5909..0ab01bb1a3d33 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/EndTxnResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/EndTxnResponse.java @@ -50,6 +50,10 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } public Errors error() { return Errors.forCode(data.errorCode()); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/EnvelopeResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/EnvelopeResponse.java index 529f616bb26fc..4f534b6721f4e 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/EnvelopeResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/EnvelopeResponse.java @@ -67,6 +67,11 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + // Not supported by the response schema + } + public static EnvelopeResponse parse(ByteBuffer buffer, short version) { return new EnvelopeResponse(new EnvelopeResponseData(new ByteBufferAccessor(buffer), version)); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ExpireDelegationTokenResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/ExpireDelegationTokenResponse.java index 163ee78d0adcd..ec43f3371b569 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/ExpireDelegationTokenResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/ExpireDelegationTokenResponse.java @@ -56,6 +56,11 @@ public ExpireDelegationTokenResponseData data() { return data; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public int throttleTimeMs() { return data.throttleTimeMs(); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/FetchResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/FetchResponse.java index a4af4ca2a2370..cd177945830d3 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/FetchResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/FetchResponse.java @@ -128,6 +128,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public int sessionId() { return data.sessionId(); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/FetchSnapshotResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/FetchSnapshotResponse.java index 7c1ce27f3da8e..d9abff66fe96e 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/FetchSnapshotResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/FetchSnapshotResponse.java @@ -61,6 +61,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public FetchSnapshotResponseData data() { return data; diff --git a/clients/src/main/java/org/apache/kafka/common/requests/FindCoordinatorResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/FindCoordinatorResponse.java index 080ba24c3bd9d..e96e8a0c0db9a 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/FindCoordinatorResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/FindCoordinatorResponse.java @@ -64,6 +64,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public boolean hasError() { return error() != Errors.NONE; } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/HeartbeatResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/HeartbeatResponse.java index eb402fcbab9f7..aebb903e967e7 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/HeartbeatResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/HeartbeatResponse.java @@ -48,6 +48,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public Errors error() { return Errors.forCode(data.errorCode()); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/IncrementalAlterConfigsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/IncrementalAlterConfigsResponse.java index b5887de9b4b75..826be30a8d3fc 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/IncrementalAlterConfigsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/IncrementalAlterConfigsResponse.java @@ -90,6 +90,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public static IncrementalAlterConfigsResponse parse(ByteBuffer buffer, short version) { return new IncrementalAlterConfigsResponse(new IncrementalAlterConfigsResponseData( new ByteBufferAccessor(buffer), version)); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/InitProducerIdResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/InitProducerIdResponse.java index f8451d7863b3f..96c7a4d400ced 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/InitProducerIdResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/InitProducerIdResponse.java @@ -48,6 +48,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { return errorCounts(Errors.forCode(data.errorCode())); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/JoinGroupResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/JoinGroupResponse.java index 336c82462a21d..5a4332efde9a9 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/JoinGroupResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/JoinGroupResponse.java @@ -47,6 +47,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public Errors error() { return Errors.forCode(data.errorCode()); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/LeaderAndIsrResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/LeaderAndIsrResponse.java index c7c04e2d99b41..0d40581d68e8a 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/LeaderAndIsrResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/LeaderAndIsrResponse.java @@ -99,6 +99,11 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + // Not supported by the response schema + } + public static LeaderAndIsrResponse parse(ByteBuffer buffer, short version) { return new LeaderAndIsrResponse(new LeaderAndIsrResponseData(new ByteBufferAccessor(buffer), version), version); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/LeaveGroupResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/LeaveGroupResponse.java index 9a59139f4e77c..d39766be68f82 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/LeaveGroupResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/LeaveGroupResponse.java @@ -82,6 +82,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public List memberResponses() { return data.members(); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsResponse.java index 270c43c0568ad..a12f85341d6a4 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsResponse.java @@ -43,6 +43,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { return errorCounts(Errors.forCode(data.errorCode())); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ListOffsetsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/ListOffsetsResponse.java index 8c4a51b542b47..53356dd93be0a 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/ListOffsetsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/ListOffsetsResponse.java @@ -64,6 +64,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public ListOffsetsResponseData data() { return data; diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ListPartitionReassignmentsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/ListPartitionReassignmentsResponse.java index 4a890e8b50cd1..cbf06d4c46624 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/ListPartitionReassignmentsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/ListPartitionReassignmentsResponse.java @@ -53,6 +53,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { return errorCounts(Errors.forCode(data.errorCode())); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ListTransactionsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/ListTransactionsResponse.java index 13ed184fc3408..f509543025b92 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/ListTransactionsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/ListTransactionsResponse.java @@ -59,4 +59,9 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/MetadataResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/MetadataResponse.java index 3696b047abad1..47cdd3f0d7e90 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/MetadataResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/MetadataResponse.java @@ -84,6 +84,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + /** * Get a map of the topics which had metadata errors * @return the map diff --git a/clients/src/main/java/org/apache/kafka/common/requests/OffsetCommitResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/OffsetCommitResponse.java index 2ed0e312983ce..713b68974a11d 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/OffsetCommitResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/OffsetCommitResponse.java @@ -107,6 +107,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public boolean shouldClientThrottle(short version) { return version >= 4; diff --git a/clients/src/main/java/org/apache/kafka/common/requests/OffsetDeleteResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/OffsetDeleteResponse.java index 79f6f4e6d3495..993a589af69fe 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/OffsetDeleteResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/OffsetDeleteResponse.java @@ -77,6 +77,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public boolean shouldClientThrottle(short version) { return version >= 0; diff --git a/clients/src/main/java/org/apache/kafka/common/requests/OffsetFetchResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/OffsetFetchResponse.java index 4e25984668da5..2d585a582ae87 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/OffsetFetchResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/OffsetFetchResponse.java @@ -245,6 +245,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public boolean hasError() { return error != Errors.NONE; } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/OffsetsForLeaderEpochResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/OffsetsForLeaderEpochResponse.java index 893d5a2af20a0..10c257c0a37cf 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/OffsetsForLeaderEpochResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/OffsetsForLeaderEpochResponse.java @@ -68,6 +68,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public static OffsetsForLeaderEpochResponse parse(ByteBuffer buffer, short version) { return new OffsetsForLeaderEpochResponse(new OffsetForLeaderEpochResponseData(new ByteBufferAccessor(buffer), version)); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ProduceResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/ProduceResponse.java index 7c9d70b8d6784..a00fdecc58645 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/ProduceResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/ProduceResponse.java @@ -116,6 +116,11 @@ public int throttleTimeMs() { return this.data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { Map errorCounts = new HashMap<>(); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/RenewDelegationTokenResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/RenewDelegationTokenResponse.java index 30708ff038c25..8ea85d74db6c7 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/RenewDelegationTokenResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/RenewDelegationTokenResponse.java @@ -53,6 +53,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public Errors error() { return Errors.forCode(data.errorCode()); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/SaslAuthenticateResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/SaslAuthenticateResponse.java index bd12d3d4ae7bb..d6ca8c170dc45 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/SaslAuthenticateResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/SaslAuthenticateResponse.java @@ -67,6 +67,11 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + // Not supported by the response schema + } + @Override public SaslAuthenticateResponseData data() { return data; diff --git a/clients/src/main/java/org/apache/kafka/common/requests/SaslHandshakeResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/SaslHandshakeResponse.java index 63c047a06196b..5097711e73787 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/SaslHandshakeResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/SaslHandshakeResponse.java @@ -57,6 +57,11 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + // Not supported by the response schema + } + @Override public SaslHandshakeResponseData data() { return data; diff --git a/clients/src/main/java/org/apache/kafka/common/requests/StopReplicaResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/StopReplicaResponse.java index 10ab153f440b9..cb66f4915d168 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/StopReplicaResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/StopReplicaResponse.java @@ -70,6 +70,11 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + // Not supported by the response schema + } + @Override public StopReplicaResponseData data() { return data; diff --git a/clients/src/main/java/org/apache/kafka/common/requests/SyncGroupResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/SyncGroupResponse.java index 822a3e78b9949..596110242902c 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/SyncGroupResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/SyncGroupResponse.java @@ -38,6 +38,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + public Errors error() { return Errors.forCode(data.errorCode()); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/TxnOffsetCommitResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/TxnOffsetCommitResponse.java index b4de54741e6d0..18244fcb17c33 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/TxnOffsetCommitResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/TxnOffsetCommitResponse.java @@ -87,6 +87,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { return errorCounts(data.topics().stream().flatMap(topic -> diff --git a/clients/src/main/java/org/apache/kafka/common/requests/UnregisterBrokerResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/UnregisterBrokerResponse.java index b508ac3ef9e54..623e6f28076fa 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/UnregisterBrokerResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/UnregisterBrokerResponse.java @@ -44,6 +44,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public Map errorCounts() { Map errorCounts = new HashMap<>(); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/UpdateFeaturesResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/UpdateFeaturesResponse.java index 26825a0c24763..567464d85dbae 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/UpdateFeaturesResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/UpdateFeaturesResponse.java @@ -63,6 +63,11 @@ public int throttleTimeMs() { return data.throttleTimeMs(); } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + @Override public String toString() { return data.toString(); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/UpdateMetadataResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/UpdateMetadataResponse.java index cc7749a47242c..d5960d7cbb923 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/UpdateMetadataResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/UpdateMetadataResponse.java @@ -47,6 +47,11 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + // Not supported by the response schema + } + public static UpdateMetadataResponse parse(ByteBuffer buffer, short version) { return new UpdateMetadataResponse(new UpdateMetadataResponseData(new ByteBufferAccessor(buffer), version)); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/VoteResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/VoteResponse.java index 51991adcf0cbb..f79c6eeb0de19 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/VoteResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/VoteResponse.java @@ -92,6 +92,11 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + // Not supported by the response schema + } + public static VoteResponse parse(ByteBuffer buffer, short version) { return new VoteResponse(new VoteResponseData(new ByteBufferAccessor(buffer), version)); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/WriteTxnMarkersResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/WriteTxnMarkersResponse.java index fd2a834d24569..a7d22e4493e67 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/WriteTxnMarkersResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/WriteTxnMarkersResponse.java @@ -108,6 +108,11 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + // Not supported by the response schema + } + @Override public Map errorCounts() { Map errorCounts = new HashMap<>(); diff --git a/core/src/main/scala/kafka/server/RequestHandlerHelper.scala b/core/src/main/scala/kafka/server/RequestHandlerHelper.scala index a1aab617b3a28..5db595986efb3 100644 --- a/core/src/main/scala/kafka/server/RequestHandlerHelper.scala +++ b/core/src/main/scala/kafka/server/RequestHandlerHelper.scala @@ -95,10 +95,13 @@ class RequestHandlerHelper( def sendForwardedResponse(request: RequestChannel.Request, response: AbstractResponse): Unit = { - // For forwarded requests, we take the throttle time from the broker that - // the request was forwarded to - val throttleTimeMs = response.throttleTimeMs() - throttle(quotas.request, request, throttleTimeMs) + // For requests forwarded to the controller, we take the maximum of the local + // request throttle and the throttle sent by the controller in the response. + val controllerThrottleTimeMs = response.throttleTimeMs() + val requestThrottleTimeMs = maybeRecordAndGetThrottleTimeMs(request) + val appliedThrottleTimeMs = math.max(controllerThrottleTimeMs, requestThrottleTimeMs) + throttle(quotas.request, request, appliedThrottleTimeMs) + response.maybeSetThrottleTimeMs(appliedThrottleTimeMs) requestChannel.sendResponse(request, response, None) } diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index aefb1a5348d77..d0e5687f67a62 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -23,7 +23,6 @@ import java.util import java.util.Arrays.asList import java.util.concurrent.TimeUnit import java.util.{Collections, Optional, Properties, Random} - import kafka.api.LeaderAndIsr import kafka.cluster.Broker import kafka.controller.{ControllerContext, KafkaController} @@ -84,7 +83,7 @@ import org.apache.kafka.server.authorizer.{Action, AuthorizationResult, Authoriz import org.junit.jupiter.api.Assertions._ import org.junit.jupiter.api.{AfterEach, Test} import org.junit.jupiter.params.ParameterizedTest -import org.junit.jupiter.params.provider.ValueSource +import org.junit.jupiter.params.provider.{CsvSource, ValueSource} import org.mockito.ArgumentMatchers.{any, anyBoolean, anyDouble, anyInt, anyLong, anyShort, anyString, argThat, isNotNull} import org.mockito.Mockito.{mock, reset, times, verify, when} import org.mockito.{ArgumentCaptor, ArgumentMatchers, Mockito} @@ -92,6 +91,7 @@ import org.mockito.{ArgumentCaptor, ArgumentMatchers, Mockito} import scala.collection.{Map, Seq, mutable} import scala.jdk.CollectionConverters._ import org.apache.kafka.common.message.CreatePartitionsRequestData.CreatePartitionsTopic +import org.apache.kafka.common.message.CreateTopicsResponseData.CreatableTopicResult import org.apache.kafka.server.common.MetadataVersion import org.apache.kafka.server.common.MetadataVersion.{IBP_0_10_2_IV0, IBP_2_2_IV1} @@ -772,6 +772,57 @@ class KafkaApisTest { testForwardableApi(ApiKeys.CREATE_TOPICS, requestBuilder) } + @ParameterizedTest + @CsvSource(value = Array("0,1500", "1500,0", "3000,1000")) + def testKRaftControllerThrottleTimeEnforced( + controllerThrottleTimeMs: Int, + requestThrottleTimeMs: Int + ): Unit = { + metadataCache = MetadataCache.kRaftMetadataCache(brokerId) + + val topicToCreate = new CreatableTopic() + .setName("topic") + .setNumPartitions(1) + .setReplicationFactor(1.toShort) + + val requestData = new CreateTopicsRequestData() + requestData.topics().add(topicToCreate) + + val requestBuilder = new CreateTopicsRequest.Builder(requestData).build() + val request = buildRequest(requestBuilder) + + val kafkaApis = createKafkaApis(enableForwarding = true, raftSupport = true) + val forwardCallback: ArgumentCaptor[Option[AbstractResponse] => Unit] = + ArgumentCaptor.forClass(classOf[Option[AbstractResponse] => Unit]) + + when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(request, time.milliseconds())) + .thenReturn(requestThrottleTimeMs) + + kafkaApis.handle(request, RequestLocal.withThreadConfinedCaching) + + verify(forwardingManager).forwardRequest( + ArgumentMatchers.eq(request), + forwardCallback.capture() + ) + + val responseData = new CreateTopicsResponseData() + .setThrottleTimeMs(controllerThrottleTimeMs) + responseData.topics().add(new CreatableTopicResult() + .setErrorCode(Errors.THROTTLING_QUOTA_EXCEEDED.code)) + + forwardCallback.getValue.apply(Some(new CreateTopicsResponse(responseData))) + + val expectedThrottleTimeMs = math.max(controllerThrottleTimeMs, requestThrottleTimeMs) + + verify(clientRequestQuotaManager).throttle( + ArgumentMatchers.eq(request), + any[ThrottleCallback](), + ArgumentMatchers.eq(expectedThrottleTimeMs) + ) + + assertEquals(expectedThrottleTimeMs, responseData.throttleTimeMs) + } + @Test def testCreatePartitionsAuthorization(): Unit = { val authorizer: Authorizer = mock(classOf[Authorizer])