This is an automated email from the ASF dual-hosted git repository.
dajac pushed a commit to branch 4.2
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/4.2 by this push:
new 61a7601d16a KAFKA-20662: Throw `INCONSISTENT_GROUP_PROTOCOL` for
malformed request during online migration (#22674)
61a7601d16a is described below
commit 61a7601d16aee7a313303f3158514b384a8fef7a
Author: Dongnuo Lyu <[email protected]>
AuthorDate: Fri Jun 26 03:08:55 2026 -0400
KAFKA-20662: Throw `INCONSISTENT_GROUP_PROTOCOL` for malformed request
during online migration (#22674)
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 23e50122a9f..902078d6e81 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
@@ -1905,7 +1905,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 c2fe7aa49e6..93f277161ea 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
@@ -12097,7 +12097,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());
}