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

  1. Upgrade kafka-clients to match the broker version.
  2. Restrict the API versions negotiated with the broker to those the client supports.
  3. 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

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


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)