This is an automated email from the ASF dual-hosted git repository.

frankvicky pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git


The following commit(s) were added to refs/heads/trunk by this push:
     new 239a3e49909 KAFKA-18157: Consider UnsupportedVersionException child 
class to represent the case of unsupported fields (#22405)
239a3e49909 is described below

commit 239a3e49909d70dfcf2bc00fab93fa7e244333fd
Author: Alan Lau <[email protected]>
AuthorDate: Tue Aug 4 11:47:26 2026 -0400

    KAFKA-18157: Consider UnsupportedVersionException child class to represent 
the case of unsupported fields (#22405)
    
    Jira: https://issues.apache.org/jira/browse/KAFKA-18157
    
    This is a rebase of #18072 onto current trunk, plus the additional scope
    requested
    
    Add a new exception class, `UnsupportedProtocolFieldException`, as a
    subclass of `UnsupportedVersionException`, to provide a more specific
    message for unsupported fields.
    
      `ConsumerHeartbeatRequestManager` uses `instanceof
    UnsupportedProtocolFieldException`    to detect this case instead of
    comparing message strings. Any future field-validation    errors added
    to `ConsumerGroupHeartbeatRequest.Builder` will be caught by the same
    `instanceof` check automatically, but the error message propagated to
    the user will need to    be updated accordingly.
    
    Reviewers: TengYao Chi <[email protected]>
---
 .../internals/ConsumerHeartbeatRequestManager.java | 10 +++--
 .../UnsupportedProtocolFieldException.java         | 40 ++++++++++++++++++
 .../requests/ConsumerGroupHeartbeatRequest.java    |  4 +-
 .../kafka/common/requests/CreateAclsRequest.java   | 14 +++++--
 .../kafka/common/requests/CreateTopicsRequest.java | 12 +++---
 .../kafka/common/requests/DeleteAclsRequest.java   |  8 ++--
 .../kafka/common/requests/DescribeAclsRequest.java |  7 ++--
 .../kafka/common/requests/ElectLeadersRequest.java |  4 +-
 .../common/requests/FindCoordinatorRequest.java    |  5 ++-
 .../kafka/common/requests/HeartbeatRequest.java    |  5 +--
 .../kafka/common/requests/JoinGroupRequest.java    |  5 +--
 .../kafka/common/requests/ListGroupsRequest.java   |  9 ++--
 .../common/requests/ListTransactionsRequest.java   |  8 ++--
 .../kafka/common/requests/MetadataRequest.java     | 10 ++---
 .../kafka/common/requests/OffsetCommitRequest.java |  4 +-
 .../kafka/common/requests/OffsetFetchRequest.java  |  4 +-
 .../kafka/common/requests/SyncGroupRequest.java    |  5 +--
 .../apache/kafka/clients/NetworkClientTest.java    |  4 +-
 .../ConsumerHeartbeatRequestManagerTest.java       | 14 +++++--
 .../UnsupportedProtocolFieldExceptionTest.java     | 49 ++++++++++++++++++++++
 .../kafka/common/requests/RequestResponseTest.java |  8 ++--
 .../processor/internals/InternalTopicManager.java  | 20 ++++-----
 22 files changed, 171 insertions(+), 78 deletions(-)

diff --git 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManager.java
 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManager.java
index e73523254a1..2269c44eaac 100644
--- 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManager.java
+++ 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManager.java
@@ -22,6 +22,7 @@ import 
org.apache.kafka.clients.consumer.internals.events.BackgroundEventHandler
 import 
org.apache.kafka.clients.consumer.internals.metrics.HeartbeatMetricsManager;
 import org.apache.kafka.common.Uuid;
 import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.ConsumerGroupHeartbeatRequestData;
 import org.apache.kafka.common.metrics.Metrics;
 import org.apache.kafka.common.protocol.Errors;
@@ -40,7 +41,6 @@ import java.util.TreeSet;
 import java.util.stream.Collectors;
 
 import static 
org.apache.kafka.clients.consumer.CloseOptions.GroupMembershipOperation.REMAIN_IN_GROUP;
-import static 
org.apache.kafka.common.requests.ConsumerGroupHeartbeatRequest.REGEX_RESOLUTION_NOT_SUPPORTED_MSG;
 
 /**
  * This is the heartbeat request manager for consumer groups.
@@ -100,12 +100,14 @@ public class ConsumerHeartbeatRequestManager extends 
AbstractHeartbeatRequestMan
         String errorMessage = exception.getMessage();
         if (exception instanceof UnsupportedVersionException) {
             String message = CONSUMER_PROTOCOL_NOT_SUPPORTED_MSG;
-            if (errorMessage.equals(REGEX_RESOLUTION_NOT_SUPPORTED_MSG)) {
-                message = REGEX_RESOLUTION_NOT_SUPPORTED_MSG;
-                logger.error("{} regex resolution not supported: {}", 
heartbeatRequestName(), message);
+            if (exception instanceof UnsupportedProtocolFieldException) {
+                message = errorMessage;
+                logger.error("{} failed due to unsupported protocol field 
while sending request: {}", heartbeatRequestName(), errorMessage);
             } else {
                 logger.error("{} failed due to unsupported version while 
sending request: {}", heartbeatRequestName(), errorMessage);
             }
+            // Surface the parent type here: handleFatalFailure propagates 
this via BackgroundEvent
+            // to the user-facing API. Propagating the subclass would be a 
user-visible behavior change.
             handleFatalFailure(new UnsupportedVersionException(message, 
exception));
             errorHandled = true;
         }
diff --git 
a/clients/src/main/java/org/apache/kafka/common/internals/UnsupportedProtocolFieldException.java
 
b/clients/src/main/java/org/apache/kafka/common/internals/UnsupportedProtocolFieldException.java
new file mode 100644
index 00000000000..8a7396febcb
--- /dev/null
+++ 
b/clients/src/main/java/org/apache/kafka/common/internals/UnsupportedProtocolFieldException.java
@@ -0,0 +1,40 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.common.internals;
+
+import org.apache.kafka.common.errors.UnsupportedVersionException;
+
+/**
+ * Indicates that a request contains a field or field value that is not 
supported by the API version
+ * negotiated with the broker, and that a higher version would be required to 
use it. This is a more
+ * specific subtype of {@link UnsupportedVersionException} that lets callers 
distinguish an unsupported
+ * field from a wholly unsupported API version.
+ */
+public class UnsupportedProtocolFieldException extends 
UnsupportedVersionException {
+    private static final long serialVersionUID = 1L;
+
+    public UnsupportedProtocolFieldException(String fieldOrValue, String 
apiKeyName,
+                                             int apiVersion, int 
lowestSupportedVersion) {
+        super("The cluster does not support [" + fieldOrValue + "] in " + 
apiKeyName
+            + " API version " + apiVersion + ". Upgrade the cluster to " + 
apiKeyName
+            + " API version >= " + lowestSupportedVersion + " to enable [" + 
fieldOrValue + "].");
+    }
+
+    public UnsupportedProtocolFieldException(String message) {
+        super(message);
+    }
+}
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/ConsumerGroupHeartbeatRequest.java
 
b/clients/src/main/java/org/apache/kafka/common/requests/ConsumerGroupHeartbeatRequest.java
index 654c0772131..d888fae36df 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/ConsumerGroupHeartbeatRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/ConsumerGroupHeartbeatRequest.java
@@ -16,7 +16,7 @@
  */
 package org.apache.kafka.common.requests;
 
-import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.ConsumerGroupHeartbeatRequestData;
 import org.apache.kafka.common.message.ConsumerGroupHeartbeatResponseData;
 import org.apache.kafka.common.protocol.ApiKeys;
@@ -63,7 +63,7 @@ public class ConsumerGroupHeartbeatRequest extends 
AbstractRequest {
         @Override
         public ConsumerGroupHeartbeatRequest build(short version) {
             if (version == 0 && data.subscribedTopicRegex() != null) {
-                throw new 
UnsupportedVersionException(REGEX_RESOLUTION_NOT_SUPPORTED_MSG);
+                throw new 
UnsupportedProtocolFieldException(REGEX_RESOLUTION_NOT_SUPPORTED_MSG);
             }
             return new ConsumerGroupHeartbeatRequest(data, version);
         }
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/CreateAclsRequest.java 
b/clients/src/main/java/org/apache/kafka/common/requests/CreateAclsRequest.java
index 5d437033208..d83572c99bd 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/CreateAclsRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/CreateAclsRequest.java
@@ -21,7 +21,7 @@ import org.apache.kafka.common.acl.AccessControlEntry;
 import org.apache.kafka.common.acl.AclBinding;
 import org.apache.kafka.common.acl.AclOperation;
 import org.apache.kafka.common.acl.AclPermissionType;
-import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.CreateAclsRequestData;
 import org.apache.kafka.common.message.CreateAclsRequestData.AclCreation;
 import org.apache.kafka.common.message.CreateAclsResponseData;
@@ -34,6 +34,7 @@ import org.apache.kafka.common.resource.ResourceType;
 
 import java.util.Collections;
 import java.util.List;
+import java.util.stream.Collectors;
 
 public class CreateAclsRequest extends AbstractRequest {
 
@@ -90,8 +91,15 @@ public class CreateAclsRequest extends AbstractRequest {
         if (version() == 0) {
             final boolean unsupported = 
data.creations().stream().anyMatch(creation ->
                 creation.resourcePatternType() != PatternType.LITERAL.code());
-            if (unsupported)
-                throw new UnsupportedVersionException("Version 0 only supports 
literal resource pattern types");
+            if (unsupported) {
+                String unsupportedType = data.creations().stream()
+                    .map(creation -> 
PatternType.fromCode(creation.resourcePatternType()))
+                    .filter(type -> type != PatternType.LITERAL)
+                    .distinct()
+                    .map(PatternType::name)
+                    .collect(Collectors.joining(","));
+                throw new UnsupportedProtocolFieldException(unsupportedType, 
apiKey().name(), version(), 1);
+            }
         }
 
         final boolean unknown = data.creations().stream().anyMatch(creation ->
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/CreateTopicsRequest.java
 
b/clients/src/main/java/org/apache/kafka/common/requests/CreateTopicsRequest.java
index ca29ba59e36..cee0c5722f1 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/CreateTopicsRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/CreateTopicsRequest.java
@@ -16,7 +16,7 @@
  */
 package org.apache.kafka.common.requests;
 
-import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.CreateTopicsRequestData;
 import org.apache.kafka.common.message.CreateTopicsRequestData.CreatableTopic;
 import org.apache.kafka.common.message.CreateTopicsResponseData;
@@ -39,8 +39,7 @@ public class CreateTopicsRequest extends AbstractRequest {
         @Override
         public CreateTopicsRequest build(short version) {
             if (data.validateOnly() && version == 0)
-                throw new UnsupportedVersionException("validateOnly is not 
supported in version 0 of " +
-                        "CreateTopicsRequest");
+                throw new UnsupportedProtocolFieldException("validateOnly", 
apiKey().name(), version, 1);
 
             final List<String> topicsWithDefaults = data.topics()
                 .stream()
@@ -52,10 +51,9 @@ public class CreateTopicsRequest extends AbstractRequest {
                 .collect(Collectors.toList());
 
             if (!topicsWithDefaults.isEmpty() && version < 4) {
-                throw new UnsupportedVersionException("Creating topics with 
default "
-                    + "partitions/replication factor are only supported in 
CreateTopicRequest "
-                    + "version 4+. The following topics need values for 
partitions and replicas: "
-                    + topicsWithDefaults);
+                throw new UnsupportedProtocolFieldException(
+                    "default partitions/replication for topics [" + 
String.join(",", topicsWithDefaults) + "]",
+                    apiKey().name(), version, 4);
             }
 
             return new CreateTopicsRequest(data, version);
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/DeleteAclsRequest.java 
b/clients/src/main/java/org/apache/kafka/common/requests/DeleteAclsRequest.java
index bb7db5b78d8..31157e6eca5 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/DeleteAclsRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/DeleteAclsRequest.java
@@ -20,7 +20,7 @@ import org.apache.kafka.common.acl.AccessControlEntryFilter;
 import org.apache.kafka.common.acl.AclBindingFilter;
 import org.apache.kafka.common.acl.AclOperation;
 import org.apache.kafka.common.acl.AclPermissionType;
-import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.DeleteAclsRequestData;
 import org.apache.kafka.common.message.DeleteAclsRequestData.DeleteAclsFilter;
 import org.apache.kafka.common.message.DeleteAclsResponseData;
@@ -76,9 +76,9 @@ public class DeleteAclsRequest extends AbstractRequest {
                 // to LITERAL. Note that the wildcard `*` is considered 
`LITERAL` for compatibility reasons.
                 if (patternType == PatternType.ANY)
                     filter.setPatternTypeFilter(PatternType.LITERAL.code());
-                else if (patternType != PatternType.LITERAL)
-                    throw new UnsupportedVersionException("Version 0 does not 
support pattern type " +
-                            patternType + " (only LITERAL and ANY are 
supported)");
+                else if (patternType != PatternType.LITERAL) {
+                    throw new 
UnsupportedProtocolFieldException(patternType.name(), apiKey().name(), 
version(), 1);
+                }
             }
         }
 
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java
 
b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java
index 8bbeeceabbf..4a96d54ffbd 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java
@@ -20,7 +20,7 @@ import org.apache.kafka.common.acl.AccessControlEntryFilter;
 import org.apache.kafka.common.acl.AclBindingFilter;
 import org.apache.kafka.common.acl.AclOperation;
 import org.apache.kafka.common.acl.AclPermissionType;
-import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.DescribeAclsRequestData;
 import org.apache.kafka.common.message.DescribeAclsResponseData;
 import org.apache.kafka.common.protocol.ApiKeys;
@@ -75,8 +75,9 @@ public class DescribeAclsRequest extends AbstractRequest {
             // to LITERAL. Note that the wildcard `*` is considered `LITERAL` 
for compatibility reasons.
             if (patternType == PatternType.ANY)
                 data.setPatternTypeFilter(PatternType.LITERAL.code());
-            else if (patternType != PatternType.LITERAL)
-                throw new UnsupportedVersionException("Version 0 only supports 
literal resource pattern types");
+            else if (patternType != PatternType.LITERAL) {
+                throw new 
UnsupportedProtocolFieldException(patternType.name(), apiKey().name(), version, 
1);
+            }
         }
 
         if (data.patternTypeFilter() == PatternType.UNKNOWN.code()
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/ElectLeadersRequest.java
 
b/clients/src/main/java/org/apache/kafka/common/requests/ElectLeadersRequest.java
index c48098ccd4b..ed7227be1b4 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/ElectLeadersRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/ElectLeadersRequest.java
@@ -19,7 +19,7 @@ package org.apache.kafka.common.requests;
 
 import org.apache.kafka.common.ElectionType;
 import org.apache.kafka.common.TopicPartition;
-import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.ElectLeadersRequestData;
 import org.apache.kafka.common.message.ElectLeadersRequestData.TopicPartitions;
 import 
org.apache.kafka.common.message.ElectLeadersResponseData.PartitionResult;
@@ -63,7 +63,7 @@ public class ElectLeadersRequest extends AbstractRequest {
 
         private ElectLeadersRequestData toRequestData(short version) {
             if (electionType != ElectionType.PREFERRED && version == 0) {
-                throw new UnsupportedVersionException("API Version 0 only 
supports PREFERRED election type");
+                throw new 
UnsupportedProtocolFieldException(electionType.name(), apiKey().name(), 
version, 1);
             }
 
             ElectLeadersRequestData data = new ElectLeadersRequestData()
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/FindCoordinatorRequest.java
 
b/clients/src/main/java/org/apache/kafka/common/requests/FindCoordinatorRequest.java
index cda50c1732f..cb7557bb0b7 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/FindCoordinatorRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/FindCoordinatorRequest.java
@@ -19,6 +19,7 @@ package org.apache.kafka.common.requests;
 import org.apache.kafka.common.Node;
 import org.apache.kafka.common.errors.InvalidRequestException;
 import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.FindCoordinatorRequestData;
 import org.apache.kafka.common.message.FindCoordinatorResponseData;
 import org.apache.kafka.common.protocol.ApiKeys;
@@ -43,8 +44,8 @@ public class FindCoordinatorRequest extends AbstractRequest {
         @Override
         public FindCoordinatorRequest build(short version) {
             if (version < 1 && data.keyType() == 
CoordinatorType.TRANSACTION.id()) {
-                throw new UnsupportedVersionException("Cannot create a v" + 
version + " FindCoordinator request " +
-                        "because we require features supported only in 2 or 
later.");
+                throw new 
UnsupportedProtocolFieldException(CoordinatorType.TRANSACTION.name(),
+                    apiKey().name(), version, 2);
             }
             int batchedKeys = data.coordinatorKeys().size();
             if (version < MIN_BATCHED_VERSION) {
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/HeartbeatRequest.java 
b/clients/src/main/java/org/apache/kafka/common/requests/HeartbeatRequest.java
index 56c7ff564a3..9e8a746722b 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/HeartbeatRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/HeartbeatRequest.java
@@ -16,7 +16,7 @@
  */
 package org.apache.kafka.common.requests;
 
-import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.HeartbeatRequestData;
 import org.apache.kafka.common.message.HeartbeatResponseData;
 import org.apache.kafka.common.protocol.ApiKeys;
@@ -36,8 +36,7 @@ public class HeartbeatRequest extends AbstractRequest {
         @Override
         public HeartbeatRequest build(short version) {
             if (data.groupInstanceId() != null && version < 3) {
-                throw new UnsupportedVersionException("The broker heartbeat 
protocol version " +
-                        version + " does not support usage of config 
group.instance.id.");
+                throw new UnsupportedProtocolFieldException("GroupInstanceId", 
apiKey().name(), version, 3);
             }
             return new HeartbeatRequest(data, version);
         }
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/JoinGroupRequest.java 
b/clients/src/main/java/org/apache/kafka/common/requests/JoinGroupRequest.java
index 30649b27303..068b80dd094 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/JoinGroupRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/JoinGroupRequest.java
@@ -17,8 +17,8 @@
 package org.apache.kafka.common.requests;
 
 import org.apache.kafka.common.errors.InvalidConfigurationException;
-import org.apache.kafka.common.errors.UnsupportedVersionException;
 import org.apache.kafka.common.internals.Topic;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.JoinGroupRequestData;
 import org.apache.kafka.common.message.JoinGroupResponseData;
 import org.apache.kafka.common.protocol.ApiKeys;
@@ -42,8 +42,7 @@ public class JoinGroupRequest extends AbstractRequest {
         @Override
         public JoinGroupRequest build(short version) {
             if (data.groupInstanceId() != null && version < 5) {
-                throw new UnsupportedVersionException("The broker join group 
protocol version " +
-                        version + " does not support usage of config 
group.instance.id.");
+                throw new UnsupportedProtocolFieldException("GroupInstanceId", 
apiKey().name(), version, 5);
             }
             return new JoinGroupRequest(data, version);
         }
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java 
b/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java
index ccdf356aff8..adb076c329f 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java
@@ -17,7 +17,7 @@
 package org.apache.kafka.common.requests;
 
 import org.apache.kafka.common.GroupType;
-import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.ListGroupsRequestData;
 import org.apache.kafka.common.message.ListGroupsResponseData;
 import org.apache.kafka.common.protocol.ApiKeys;
@@ -48,8 +48,7 @@ public class ListGroupsRequest extends AbstractRequest {
         @Override
         public ListGroupsRequest build(short version) {
             if (!data.statesFilter().isEmpty() && version < 4) {
-                throw new UnsupportedVersionException("The broker only 
supports ListGroups " +
-                        "v" + version + ", but we need v4 or newer to request 
groups by states.");
+                throw new UnsupportedProtocolFieldException("StatesFilter", 
apiKey().name(), version, 4);
             }
             if (!data.typesFilter().isEmpty() && version < 5) {
                 // Types filter is supported by brokers with version 3.8.0 or 
later. Older brokers only support
@@ -60,9 +59,7 @@ public class ListGroupsRequest extends AbstractRequest {
                 boolean containedClassic = 
typesCopy.remove(GroupType.CLASSIC.toString());
                 boolean containedConsumer = 
typesCopy.remove(GroupType.CONSUMER.toString());
                 if (!typesCopy.isEmpty() || (!containedClassic && 
containedConsumer)) {
-                    throw new UnsupportedVersionException("The broker only 
supports ListGroups " +
-                        "v" + version + ", but we need v5 or newer to request 
groups by type. " +
-                        "Requested group types: [" + String.join(", ", 
data.typesFilter()) + "].");
+                    throw new 
UnsupportedProtocolFieldException(String.join(",", data.typesFilter()), 
apiKey().name(), version, 5);
                 }
                 return new 
ListGroupsRequest(data.duplicate().setTypesFilter(List.of()), version);
             }
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/ListTransactionsRequest.java
 
b/clients/src/main/java/org/apache/kafka/common/requests/ListTransactionsRequest.java
index 34c39625972..ef06d805b9d 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/ListTransactionsRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/ListTransactionsRequest.java
@@ -16,7 +16,7 @@
  */
 package org.apache.kafka.common.requests;
 
-import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.ListTransactionsRequestData;
 import org.apache.kafka.common.message.ListTransactionsResponseData;
 import org.apache.kafka.common.protocol.ApiKeys;
@@ -35,12 +35,10 @@ public class ListTransactionsRequest extends 
AbstractRequest {
         @Override
         public ListTransactionsRequest build(short version) {
             if (data.durationFilter() >= 0 && version < 1) {
-                throw new UnsupportedVersionException("Duration filter can be 
set only when using API version 1 or higher." +
-                        " If client is connected to an older broker, do not 
specify duration filter or set duration filter to -1.");
+                throw new UnsupportedProtocolFieldException("DurationFilter", 
apiKey().name(), version, 1);
             }
             if (data.transactionalIdPattern() != null && version < 2) {
-                throw new UnsupportedVersionException("Transactional ID 
pattern filter can be set only when using API version 2 or higher." +
-                    " If client is connected to an older broker, do not 
specify the pattern filter.");
+                throw new 
UnsupportedProtocolFieldException("TransactionalIdPattern", apiKey().name(), 
version, 2);
             }
             return new ListTransactionsRequest(data, version);
         }
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/MetadataRequest.java 
b/clients/src/main/java/org/apache/kafka/common/requests/MetadataRequest.java
index 392bf2a9105..f9015392c30 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/MetadataRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/MetadataRequest.java
@@ -18,6 +18,7 @@ package org.apache.kafka.common.requests;
 
 import org.apache.kafka.common.Uuid;
 import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.MetadataRequestData;
 import 
org.apache.kafka.common.message.MetadataRequestData.MetadataRequestTopic;
 import org.apache.kafka.common.message.MetadataResponseData;
@@ -126,16 +127,13 @@ public class MetadataRequest extends AbstractRequest {
             if (version < 1)
                 throw new UnsupportedVersionException("MetadataRequest 
versions older than 1 are not supported.");
             if (!data.allowAutoTopicCreation() && version < 4)
-                throw new UnsupportedVersionException("MetadataRequest 
versions older than 4 don't support the " +
-                        "allowAutoTopicCreation field");
+                throw new 
UnsupportedProtocolFieldException("allowAutoTopicCreation", apiKey().name(), 
version, 4);
             if (data.topics() != null) {
                 data.topics().forEach(topic -> {
                     if (topic.name() == null && version < 12)
-                        throw new UnsupportedVersionException("MetadataRequest 
version " + version +
-                                " does not support null topic names.");
+                        throw new UnsupportedProtocolFieldException("null 
topic names", apiKey().name(), version, 12);
                     if (!Uuid.ZERO_UUID.equals(topic.topicId()) && version < 
12)
-                        throw new UnsupportedVersionException("MetadataRequest 
version " + version +
-                            " does not support non-zero topic IDs.");
+                        throw new UnsupportedProtocolFieldException("non-zero 
topic IDs", apiKey().name(), version, 12);
                 });
             }
             return new MetadataRequest(data, version);
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/OffsetCommitRequest.java
 
b/clients/src/main/java/org/apache/kafka/common/requests/OffsetCommitRequest.java
index d0c8813e7fb..1a762f068dc 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/OffsetCommitRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/OffsetCommitRequest.java
@@ -19,6 +19,7 @@ package org.apache.kafka.common.requests;
 import org.apache.kafka.common.TopicPartition;
 import org.apache.kafka.common.Uuid;
 import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.OffsetCommitRequestData;
 import 
org.apache.kafka.common.message.OffsetCommitRequestData.OffsetCommitRequestTopic;
 import org.apache.kafka.common.message.OffsetCommitResponseData;
@@ -62,8 +63,7 @@ public class OffsetCommitRequest extends AbstractRequest {
         @Override
         public OffsetCommitRequest build(short version) {
             if (data.groupInstanceId() != null && version < 7) {
-                throw new UnsupportedVersionException("The broker offset 
commit api version " +
-                    version + " does not support usage of config 
group.instance.id.");
+                throw new UnsupportedProtocolFieldException("GroupInstanceId", 
apiKey().name(), version, 7);
             }
             if (version >= 10) {
                 data.topics().forEach(topic -> {
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/OffsetFetchRequest.java
 
b/clients/src/main/java/org/apache/kafka/common/requests/OffsetFetchRequest.java
index 3cb9588c4b7..507c2ef4a49 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/OffsetFetchRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/OffsetFetchRequest.java
@@ -19,6 +19,7 @@ package org.apache.kafka.common.requests;
 import org.apache.kafka.common.TopicPartition;
 import org.apache.kafka.common.Uuid;
 import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.OffsetFetchRequestData;
 import 
org.apache.kafka.common.message.OffsetFetchRequestData.OffsetFetchRequestGroup;
 import 
org.apache.kafka.common.message.OffsetFetchRequestData.OffsetFetchRequestTopic;
@@ -97,8 +98,7 @@ public class OffsetFetchRequest extends AbstractRequest {
         private void throwIfStableOffsetsUnsupported(short version) {
             if (data.requireStable() && version < 
REQUIRE_STABLE_OFFSET_MIN_VERSION) {
                 if (throwOnFetchStableOffsetsUnsupported) {
-                    throw new UnsupportedVersionException("Broker unexpectedly 
" +
-                        "doesn't support requireStable flag on version " + 
version);
+                    throw new 
UnsupportedProtocolFieldException("RequireStable", apiKey().name(), version, 7);
                 } else {
                     log.trace("Fallback the requireStable flag to false as 
broker " +
                         "only supports OffsetFetchRequest version {}. Need " +
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/SyncGroupRequest.java 
b/clients/src/main/java/org/apache/kafka/common/requests/SyncGroupRequest.java
index ceb0e5248c7..b99139fd03a 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/SyncGroupRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/SyncGroupRequest.java
@@ -16,7 +16,7 @@
  */
 package org.apache.kafka.common.requests;
 
-import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.SyncGroupRequestData;
 import org.apache.kafka.common.message.SyncGroupResponseData;
 import org.apache.kafka.common.protocol.ApiKeys;
@@ -41,8 +41,7 @@ public class SyncGroupRequest extends AbstractRequest {
         @Override
         public SyncGroupRequest build(short version) {
             if (data.groupInstanceId() != null && version < 3) {
-                throw new UnsupportedVersionException("The broker sync group 
protocol version " +
-                        version + " does not support usage of config 
group.instance.id.");
+                throw new UnsupportedProtocolFieldException("GroupInstanceId", 
apiKey().name(), version, 3);
             }
             return new SyncGroupRequest(data, version);
         }
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java 
b/clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java
index 0111c6baada..fa1b4c4176e 100644
--- a/clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java
+++ b/clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java
@@ -21,8 +21,8 @@ import org.apache.kafka.common.KafkaException;
 import org.apache.kafka.common.Node;
 import org.apache.kafka.common.errors.AuthenticationException;
 import org.apache.kafka.common.errors.RebootstrapRequiredException;
-import org.apache.kafka.common.errors.UnsupportedVersionException;
 import org.apache.kafka.common.internals.ClusterResourceListeners;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.ApiMessageType;
 import org.apache.kafka.common.message.ApiVersionsResponseData;
 import org.apache.kafka.common.message.ApiVersionsResponseData.ApiVersion;
@@ -259,7 +259,7 @@ public class NetworkClientTest {
         // disabling auto topic creation for versions less than 4 is not 
supported
         MetadataRequest.Builder builder = new MetadataRequest.Builder(topics, 
false, (short) 3);
         client.sendInternalMetadataRequest(builder, node.idString(), 
time.milliseconds());
-        assertEquals(UnsupportedVersionException.class, 
metadataUpdater.getAndClearFailure().getClass());
+        assertEquals(UnsupportedProtocolFieldException.class, 
metadataUpdater.getAndClearFailure().getClass());
     }
 
     @Test
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java
index b873571e00c..5d3f6a7b285 100644
--- 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java
@@ -32,6 +32,7 @@ import org.apache.kafka.common.errors.AuthenticationException;
 import org.apache.kafka.common.errors.DisconnectException;
 import org.apache.kafka.common.errors.TimeoutException;
 import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.ConsumerGroupHeartbeatRequestData;
 import org.apache.kafka.common.message.ConsumerGroupHeartbeatResponseData;
 import org.apache.kafka.common.metrics.Metrics;
@@ -542,9 +543,9 @@ public class ConsumerHeartbeatRequestManagerTest
      * REGEX_RESOLUTION_NOT_SUPPORTED_MSG only generated on the client side.
      */
     @ParameterizedTest
-    @ValueSource(strings = {CONSUMER_PROTOCOL_NOT_SUPPORTED_MSG, 
REGEX_RESOLUTION_NOT_SUPPORTED_MSG})
-    public void testUnsupportedVersionFromClient(String errorMsg) {
-        mockResponseWithException(new UnsupportedVersionException(errorMsg), 
false);
+    @MethodSource("unsupportedVersionFromClientCases")
+    public void testUnsupportedVersionFromClient(UnsupportedVersionException 
thrown, String errorMsg) {
+        mockResponseWithException(thrown, false);
         ArgumentCaptor<ErrorEvent> errorEventArgumentCaptor = 
ArgumentCaptor.forClass(ErrorEvent.class);
         verify(backgroundEventHandler).add(errorEventArgumentCaptor.capture());
         ErrorEvent errorEvent = errorEventArgumentCaptor.getValue();
@@ -553,6 +554,13 @@ public class ConsumerHeartbeatRequestManagerTest
         clearInvocations(backgroundEventHandler);
     }
 
+    private static Stream<Arguments> unsupportedVersionFromClientCases() {
+        return Stream.of(
+            Arguments.of(new 
UnsupportedVersionException(CONSUMER_PROTOCOL_NOT_SUPPORTED_MSG), 
CONSUMER_PROTOCOL_NOT_SUPPORTED_MSG),
+            Arguments.of(new 
UnsupportedProtocolFieldException(REGEX_RESOLUTION_NOT_SUPPORTED_MSG), 
REGEX_RESOLUTION_NOT_SUPPORTED_MSG)
+        );
+    }
+
     private void mockResponseWithException(UnsupportedVersionException 
exception, boolean isFromBroker) {
         time.sleep(DEFAULT_HEARTBEAT_INTERVAL_MS);
         NetworkClientDelegate.PollResult result = 
heartbeatRequestManager.poll(time.milliseconds());
diff --git 
a/clients/src/test/java/org/apache/kafka/common/internals/UnsupportedProtocolFieldExceptionTest.java
 
b/clients/src/test/java/org/apache/kafka/common/internals/UnsupportedProtocolFieldExceptionTest.java
new file mode 100644
index 00000000000..025b41fb13a
--- /dev/null
+++ 
b/clients/src/test/java/org/apache/kafka/common/internals/UnsupportedProtocolFieldExceptionTest.java
@@ -0,0 +1,49 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.common.internals;
+
+import org.apache.kafka.common.errors.UnsupportedVersionException;
+
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+
+public class UnsupportedProtocolFieldExceptionTest {
+
+    @Test
+    public void testFieldConstructorFormatsMessage() {
+        UnsupportedProtocolFieldException exception =
+            new UnsupportedProtocolFieldException("validateOnly", 
"CREATE_TOPICS", 0, 1);
+        assertEquals("The cluster does not support [validateOnly] in 
CREATE_TOPICS API version 0. "
+                + "Upgrade the cluster to CREATE_TOPICS API version >= 1 to 
enable [validateOnly].",
+            exception.getMessage());
+    }
+
+    @Test
+    public void testMessageConstructorPassesMessageThrough() {
+        UnsupportedProtocolFieldException exception =
+            new UnsupportedProtocolFieldException("some custom message");
+        assertEquals("some custom message", exception.getMessage());
+    }
+
+    @Test
+    public void testIsUnsupportedVersionException() {
+        assertInstanceOf(UnsupportedVersionException.class,
+            new UnsupportedProtocolFieldException("field", "SOME_API", 0, 1));
+    }
+}
diff --git 
a/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java
 
b/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java
index 5b1f4e76c2d..bd3773b417c 100644
--- 
a/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java
@@ -36,6 +36,7 @@ import 
org.apache.kafka.common.errors.NotEnoughReplicasException;
 import org.apache.kafka.common.errors.SecurityDisabledException;
 import org.apache.kafka.common.errors.UnknownServerException;
 import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.message.AddOffsetsToTxnRequestData;
 import org.apache.kafka.common.message.AddOffsetsToTxnResponseData;
 import 
org.apache.kafka.common.message.AddPartitionsToTxnRequestData.AddPartitionsToTxnTopic;
@@ -745,8 +746,8 @@ public class RequestResponseTest {
 
     @Test
     public void testCreateTopicRequestV3FailsIfNoPartitionsOrReplicas() {
-        final UnsupportedVersionException exception = assertThrows(
-            UnsupportedVersionException.class, () -> {
+        final UnsupportedProtocolFieldException exception = assertThrows(
+            UnsupportedProtocolFieldException.class, () -> {
                 CreateTopicsRequestData data = new CreateTopicsRequestData()
                     .setTimeoutMs(123)
                     .setValidateOnly(false);
@@ -761,8 +762,7 @@ public class RequestResponseTest {
 
                 new Builder(data).build((short) 3);
             });
-        assertTrue(exception.getMessage().contains("supported in 
CreateTopicRequest version 4+"));
-        assertTrue(exception.getMessage().contains("[foo, bar]"));
+        assertTrue(exception.getMessage().contains("does not support [default 
partitions/replication for topics [foo,bar]] in CREATE_TOPICS API version 3"));
     }
 
     @Test
diff --git 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopicManager.java
 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopicManager.java
index cfb376f0d14..8b79a7add96 100644
--- 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopicManager.java
+++ 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopicManager.java
@@ -36,7 +36,7 @@ import 
org.apache.kafka.common.errors.LeaderNotAvailableException;
 import org.apache.kafka.common.errors.TimeoutException;
 import org.apache.kafka.common.errors.TopicExistsException;
 import org.apache.kafka.common.errors.UnknownTopicOrPartitionException;
-import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.internals.UnsupportedProtocolFieldException;
 import org.apache.kafka.common.serialization.ByteArrayDeserializer;
 import org.apache.kafka.common.utils.Time;
 import org.apache.kafka.common.utils.Utils;
@@ -575,17 +575,13 @@ public class InternalTopicManager {
                     log.error("Unexpected error during topic creation for 
{}.\n" +
                             "Error message was: {}", topicName, 
cause.toString());
 
-                    if (cause instanceof UnsupportedVersionException) {
-                        final String errorMessage = cause.getMessage();
-                        if (errorMessage != null &&
-                                errorMessage.startsWith("Creating topics with 
default partitions/replication factor are only supported in CreateTopicRequest 
version 4+")) {
-
-                            throw new StreamsException(String.format(
-                                    "Could not create topic %s, because 
brokers don't support configuration replication.factor=-1."
-                                            + " You can change the 
replication.factor config or upgrade your brokers to version 2.4 or newer to 
avoid this error.",
-                                    topicName)
-                            );
-                        }
+                    if (cause instanceof UnsupportedProtocolFieldException) {
+                        // An older broker rejected a field we rely on (e.g. 
the default
+                        // replication.factor=-1, which requires CreateTopics 
request version 4+).
+                        throw new StreamsException(String.format(
+                                "Could not create topic %s, because brokers 
don't support configuration replication.factor=-1."
+                                        + " You can change the 
replication.factor config or upgrade your brokers to version 2.4 or newer to 
avoid this error.",
+                                topicName), cause);
                     } else if (cause instanceof TimeoutException) {
                         log.error("Creating topic {} timed out.\n" +
                                 "Error message was: {}", topicName, 
cause.toString());

Reply via email to