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());