junrao commented on code in PR #22989:
URL: https://github.com/apache/kafka/pull/22989#discussion_r3754292286
##########
clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java:
##########
@@ -306,6 +344,10 @@ public static AbstractResponse parseResponse(ApiKeys
apiKey, Readable readable,
* quota violation, sends out responses before throttling.
*/
public boolean shouldClientThrottle(short version) {
+ Short firstPostKip219Version = POST_KIP_219_VERSION.get(apiKey);
Review Comment:
Could we make this method final to avoid accidental override in subclasses
in the future?
##########
clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java:
##########
@@ -367,50 +361,12 @@ public void testSerialization() {
@Test
public void testClientThrottlesResponsesWithThrottleTime() {
- Map<ApiKeys, Short> postKip219Version = Map.ofEntries(
- Map.entry(ApiKeys.PRODUCE, (short) 6),
- Map.entry(ApiKeys.FETCH, (short) 8),
- Map.entry(ApiKeys.LIST_OFFSETS, (short) 3),
- Map.entry(ApiKeys.METADATA, (short) 6),
- Map.entry(ApiKeys.OFFSET_COMMIT, (short) 4),
- Map.entry(ApiKeys.OFFSET_FETCH, (short) 4),
- Map.entry(ApiKeys.FIND_COORDINATOR, (short) 2),
- Map.entry(ApiKeys.JOIN_GROUP, (short) 3),
- Map.entry(ApiKeys.HEARTBEAT, (short) 2),
- Map.entry(ApiKeys.LEAVE_GROUP, (short) 2),
- Map.entry(ApiKeys.SYNC_GROUP, (short) 2),
- Map.entry(ApiKeys.DESCRIBE_GROUPS, (short) 2),
- Map.entry(ApiKeys.LIST_GROUPS, (short) 2),
- Map.entry(ApiKeys.API_VERSIONS, (short) 2),
- Map.entry(ApiKeys.CREATE_TOPICS, (short) 3),
- Map.entry(ApiKeys.DELETE_TOPICS, (short) 2),
- Map.entry(ApiKeys.DELETE_RECORDS, (short) 1),
- Map.entry(ApiKeys.INIT_PRODUCER_ID, (short) 1),
- Map.entry(ApiKeys.ADD_PARTITIONS_TO_TXN, (short) 1),
- Map.entry(ApiKeys.ADD_OFFSETS_TO_TXN, (short) 1),
- Map.entry(ApiKeys.END_TXN, (short) 1),
- Map.entry(ApiKeys.TXN_OFFSET_COMMIT, (short) 1),
- Map.entry(ApiKeys.DESCRIBE_ACLS, (short) 1),
- Map.entry(ApiKeys.CREATE_ACLS, (short) 1),
- Map.entry(ApiKeys.DELETE_ACLS, (short) 1),
- Map.entry(ApiKeys.DESCRIBE_CONFIGS, (short) 2),
- Map.entry(ApiKeys.ALTER_CONFIGS, (short) 1),
- Map.entry(ApiKeys.ALTER_REPLICA_LOG_DIRS, (short) 1),
- Map.entry(ApiKeys.DESCRIBE_LOG_DIRS, (short) 1),
- Map.entry(ApiKeys.CREATE_PARTITIONS, (short) 1),
- Map.entry(ApiKeys.CREATE_DELEGATION_TOKEN, (short) 1),
- Map.entry(ApiKeys.RENEW_DELEGATION_TOKEN, (short) 1),
- Map.entry(ApiKeys.EXPIRE_DELEGATION_TOKEN, (short) 1),
- Map.entry(ApiKeys.DESCRIBE_DELEGATION_TOKEN, (short) 1),
- Map.entry(ApiKeys.DELETE_GROUPS, (short) 1)
- );
-
for (ApiKeys apiKey : ApiKeys.values()) {
for (short version : apiKey.allVersions()) {
boolean responseHasThrottleTime =
apiKey.messageType.responseSchemas()[version].get("throttle_time_ms") != null;
- short firstPostKip219Version =
postKip219Version.getOrDefault(apiKey, apiKey.oldestVersion());
- boolean shouldClientThrottle = responseHasThrottleTime &&
+ Short firstPostKip219Version =
AbstractResponse.POST_KIP_219_VERSION.get(apiKey);
Review Comment:
This is now identical to the implementation. So the test doesn't seem useful.
Claude suggested the following.
1. We can add the following to test to catch any boundary set too low.
```
@Test
public void testKip219BoundariesHaveThrottleTimeField() {
AbstractResponse.POST_KIP_219_VERSION.forEach((apiKey, boundary) ->
assertNotNull(apiKey.messageType.responseSchemas()[boundary].get("throttle_time_ms"),
apiKey + " KIP-219 boundary " + boundary + " precedes the
version that added throttle_time_ms"));
}
```
2. OffsetCommitResponseTest explicitly asserts version >= 4. We can add the
same for the 3 most common type of requests producer, fetch and metadata.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]