apache/kafka · error · IllegalArgumentException
Unknown rebalance protocol id
Error message
Unknown rebalance protocol id: {id} What it means
Thrown by ConsumerPartitionAssignor.RebalanceProtocol.forId(byte) when the id is not 0 (EAGER) or 1 (COOPERATIVE). The enum is used when decoding the rebalance protocol carried over the wire (e.g. in SyncGroup/Heartbeat responses). An unknown id means the broker/assignor produced a protocol value this client does not recognize.
Solutions
- Upgrade kafka-clients to match the broker version.
- Restrict the API versions negotiated with the broker to those the client supports.
- Catch IllegalArgumentException and treat the response as unsupported for this client.
Example fix
// before
RebalanceProtocol p = RebalanceProtocol.forId(rawId);
// after
RebalanceProtocol p;
try {
p = RebalanceProtocol.forId(rawId);
} catch (IllegalArgumentException e) {
throw new IllegalStateException("Unsupported rebalance protocol id " + rawId + "; upgrade kafka-clients", e);
} Defensive patterns
Strategy: try-catch
Validate before calling
// Pre-validate against known ids
if (rawId != 0 && rawId != 1) {
throw new IllegalArgumentException("Unsupported rebalance protocol id " + rawId);
}
ConsumerPartitionAssignor.RebalanceProtocol.forId(rawId); Type guard
static boolean isKnownRebalanceProtocolId(byte id) {
return id == 0 || id == 1;
} Try / catch
try {
protocol = ConsumerPartitionAssignor.RebalanceProtocol.forId(rawId);
} catch (IllegalArgumentException e) {
throw new IllegalStateException("Unsupported rebalance protocol id " + rawId + "; upgrade kafka-clients", e);
} Prevention
- Keep client and broker versions aligned.
- Restrict negotiated API versions on the broker to those the client supports.
- Surface unknown-wire-id failures as version-mismatch diagnostics.
When it happens
Trigger: Calling RebalanceProtocol.forId with a byte other than 0 or 1; deserializing a protocol message from a newer broker that added a new rebalance protocol.
Common situations: Client older than broker; experimental protocol extension; corrupt wire data.
Related errors
- Unknown acknowledge type id
- Unknown topology description status id
- Cannot add records for a partition that is not assigned to…
- Cannot lose partitions that are not currently assigned
- Global store must be composed of a source and a processor…
AI-assisted analysis of apache/kafka@996fb4585a (2026-08-11).
Data as JSON: /api/errors/5fa0f79121496a39.
Report an issue: GitHub.
Appendix: source
Thrown at clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerPartitionAssignor.java:409
public byte id() {
return id;
}
/**
* Returns the rebalance protocol for the given identifier.
*
* @param id The identifier for the rebalance protocol
* @return The corresponding rebalance protocol
* @throws IllegalArgumentException If the ID is not recognized
*/
public static RebalanceProtocol forId(byte id) {
switch (id) {
case 0:
return EAGER;
case 1:
return COOPERATIVE;
default:
throw new IllegalArgumentException("Unknown rebalance protocol id: " + id);
}
}
}
/**
* Get a list of configured instances of {@link org.apache.kafka.clients.consumer.ConsumerPartitionAssignor}
* based on the class names/types specified by {@link org.apache.kafka.clients.consumer.ConsumerConfig#PARTITION_ASSIGNMENT_STRATEGY_CONFIG}
*/
static List<ConsumerPartitionAssignor> getAssignorInstances(List<String> assignorClasses, Map<String, Object> configs) {
List<ConsumerPartitionAssignor> assignors = new ArrayList<>();
// a map to store assignor name -> assignor class name
Map<String, String> assignorNameMap = new HashMap<>();
for (Object klass : assignorClasses) {
// first try to get the class if passed in as a string
if (klass instanceof String) {
try {
klass = Utils.loadClass((String) klass, Object.class);View on GitHub (pinned to 996fb4585a)