This is an automated email from the ASF dual-hosted git repository.
dajac pushed a commit to branch 4.3
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/4.3 by this push:
new c7af7e5aca2 KAFKA-20662: Throw `INCONSISTENT_GROUP_PROTOCOL` for
malformed request during online migration (#22580)
c7af7e5aca2 is described below
commit c7af7e5aca2ee5ff8c994bf592a3d4db0c126466
Author: Dongnuo Lyu <[email protected]>
AuthorDate: Tue Jun 23 02:59:29 2026 -0400
KAFKA-20662: Throw `INCONSISTENT_GROUP_PROTOCOL` for malformed request
during online migration (#22580)
When the subscription deserialization fails, we throw an
IllegalStateException
```
private static ConsumerProtocolSubscription deserializeSubscription(
JoinGroupRequestProtocolCollection protocols
) {
try {
return ConsumerProtocol.deserializeConsumerProtocolSubscription(
ByteBuffer.wrap(protocols.iterator().next().metadata()));
}
catch (SchemaException e) {
throw new IllegalStateException("Malformed embedded consumer protocol
in subscription deserialization.");
}
}
```
which will be translated to UNKNOWN_SERVER_ERROR. This is a bit
confusing to the client. We throw a fatal `INCONSISTENT_GROUP_PROTOCOL`
here instead.
Reviewers: David Jacot <[email protected]>
---
.../java/org/apache/kafka/coordinator/group/GroupMetadataManager.java | 2 +-
.../org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java | 2 +-
2 files changed, 2 insertions(+), 2 deletions(-)
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
index e91a739fca3..b4373210d25 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
@@ -1945,7 +1945,7 @@ public class GroupMetadataManager {
ByteBuffer.wrap(protocols.iterator().next().metadata())
);
} catch (SchemaException e) {
- throw new IllegalStateException("Malformed embedded consumer
protocol in subscription deserialization.");
+ throw Errors.INCONSISTENT_GROUP_PROTOCOL.exception("Malformed
embedded consumer protocol in subscription deserialization.");
}
}
diff --git
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
index dc83f47a31b..c8f9a0d1202 100644
---
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
+++
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
@@ -12298,7 +12298,7 @@ public class GroupMetadataManagerTest {
.setSessionTimeoutMs(5000)
.setRebalanceTimeoutMs(45000);
- IllegalStateException ex = assertThrows(IllegalStateException.class,
+ InconsistentGroupProtocolException ex =
assertThrows(InconsistentGroupProtocolException.class,
() -> context.sendClassicGroupJoin(joinRequest));
assertEquals("Malformed embedded consumer protocol in subscription
deserialization.", ex.getMessage());
}