This is an automated email from the ASF dual-hosted git repository.
junrao 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 c274a7348fb KAFKA-20395: Support unregistering controllers (#22191)
c274a7348fb is described below
commit c274a7348fb66d1bde99566bdb072fbb52b03dfd
Author: Kevin Wu <[email protected]>
AuthorDate: Mon Aug 3 13:33:46 2026 -0500
KAFKA-20395: Support unregistering controllers (#22191)
### What Changed
- Introduce `UnregisterControllerRequest` and
`UnregisterControllerResponse` RPC schemas
- Introduce `Admin#unregisterController` interface, which the
`kafka-cluster` and `kafka-metadata-quorum` tools call via the
`AdminClient` to unregister a controller
- Add controller-side handling for `UnregisterControllerRequest`
- Introduce `UnregisterControllerRecord` metadata record schema + new
MetadataVersion `IBP_4_4_IV2` to support this feature
- Add handling for replaying the `UnregisterControllerRecord` on the
active controller and in the `Cluster/MetadataDelta` to remove
unregistered controllers from the `MetadataImage`
- Add the `kafka-cluster unregister-controller` CLI command
- Add the `--unregister` flag to the `kafka-metadata-quorum
remove-controller` command
### Testing
- Unit testing for CLI changes, controller-side changes, and AdminClient
changes
- Integration tests in `KRaftClusterTest`
Reviewers: Jun Rao <[email protected]>, Chetan Koneru
(github:ChetanKoneru), Paolo Patierno <[email protected]>
---
.../java/org/apache/kafka/clients/admin/Admin.java | 42 +++++
.../kafka/clients/admin/ForwardingAdmin.java | 5 +
.../kafka/clients/admin/KafkaAdminClient.java | 45 +++++
.../clients/admin/UnregisterControllerOptions.java | 27 +++
.../clients/admin/UnregisterControllerResult.java | 42 +++++
.../errors/ControllerIdNotRegisteredException.java | 32 ++++
.../org/apache/kafka/common/protocol/ApiKeys.java | 3 +-
.../org/apache/kafka/common/protocol/Errors.java | 4 +-
.../kafka/common/requests/AbstractRequest.java | 2 +
.../kafka/common/requests/AbstractResponse.java | 2 +
.../requests/UnregisterControllerRequest.java | 66 ++++++++
.../requests/UnregisterControllerResponse.java | 63 +++++++
.../message/UnregisterControllerRequest.json | 27 +++
.../message/UnregisterControllerResponse.json | 30 ++++
.../kafka/clients/admin/KafkaAdminClientTest.java | 184 +++++++++++++--------
.../kafka/common/requests/RequestResponseTest.java | 39 ++++-
.../kafka/clients/admin/MockAdminClient.java | 11 ++
.../main/scala/kafka/server/ControllerApis.scala | 22 +++
core/src/main/scala/kafka/server/KafkaApis.scala | 1 +
.../unit/kafka/server/ControllerApisTest.scala | 9 +
.../server/ControllerRegistrationManagerTest.scala | 122 +++++++++++++-
.../scala/unit/kafka/server/RequestQuotaTest.scala | 3 +
docs/getting-started/upgrade.md | 1 +
docs/operations/kraft.md | 20 +++
docs/security/authorization-and-acls.md | 17 ++
.../kafka/controller/ClusterControlManager.java | 29 ++++
.../org/apache/kafka/controller/Controller.java | 14 ++
.../apache/kafka/controller/QuorumController.java | 19 +++
.../java/org/apache/kafka/image/ClusterDelta.java | 5 +
.../java/org/apache/kafka/image/MetadataDelta.java | 8 +
.../metadata/UnregisterControllerRecord.json | 26 +++
.../controller/ClusterControlManagerTest.java | 83 ++++++++++
.../org/apache/kafka/metadata/RecordTestUtils.java | 5 +
.../kafka/server/common/MetadataVersion.java | 9 +-
.../apache/kafka/network/RequestConvertToJson.java | 8 +
.../org/apache/kafka/server/KRaftClusterTest.java | 71 ++++++++
.../TestingMetricsInterceptingAdminClient.java | 7 +
.../apache/kafka/common/test/api/ClusterTest.java | 2 +-
.../apache/kafka/common/test/MockController.java | 8 +
.../java/org/apache/kafka/tools/ClusterTool.java | 29 +++-
.../apache/kafka/tools/MetadataQuorumCommand.java | 62 ++++++-
.../org/apache/kafka/tools/ClusterToolTest.java | 20 +++
.../kafka/tools/MetadataQuorumCommandUnitTest.java | 118 ++++++++++++-
43 files changed, 1251 insertions(+), 91 deletions(-)
diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java
b/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java
index 942b71d6433..4cb40c66980 100644
--- a/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java
+++ b/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java
@@ -1682,6 +1682,48 @@ public interface Admin extends AutoCloseable {
@InterfaceStability.Unstable
UnregisterBrokerResult unregisterBroker(int brokerId,
UnregisterBrokerOptions options);
+ /**
+ * Unregister a controller.
+ *
+ * This is a convenience method for {@link #unregisterController(int,
UnregisterControllerOptions)}
+ *
+ * @param controllerId the controller id to unregister.
+ *
+ * @return the {@link UnregisterControllerResult} containing the result
+ */
+ @InterfaceStability.Unstable
+ default UnregisterControllerResult unregisterController(int controllerId) {
+ return unregisterController(controllerId, new
UnregisterControllerOptions());
+ }
+
+ /**
+ * Unregister a controller.
+ *
+ * The following exceptions can be anticipated when calling {@code get()}
on the future from the
+ * returned {@link UnregisterControllerResult}:
+ * <ul>
+ * <li>{@link org.apache.kafka.common.errors.TimeoutException}
+ * If the request timed out before the unregister operation could
finish.</li>
+ * <li>{@link org.apache.kafka.common.errors.UnsupportedVersionException}
+ * If the software is too old to support the unregistration API.</li>
+ * <li>{@link
org.apache.kafka.common.errors.ControllerIdNotRegisteredException}
+ * If the requested controller id is not currently registered.</li>
+ * <li>{@link org.apache.kafka.common.errors.NotControllerException}
+ * If the request does not arrive at the active controller.</li>
+ * <li>{@link org.apache.kafka.common.errors.InvalidRequestException}
+ * If the request tries to unregister the current active controller
id.</li>
+ *
+ * </ul>
+ * <p>
+ *
+ * @param controllerId the controller id to unregister.
+ * @param options the options to use.
+ *
+ * @return the {@link UnregisterControllerResult} containing the result
+ */
+ @InterfaceStability.Unstable
+ UnregisterControllerResult unregisterController(int controllerId,
UnregisterControllerOptions options);
+
/**
* Describe producer state on a set of topic partitions. See
* {@link #describeProducers(Collection, DescribeProducersOptions)} for
more details.
diff --git
a/clients/src/main/java/org/apache/kafka/clients/admin/ForwardingAdmin.java
b/clients/src/main/java/org/apache/kafka/clients/admin/ForwardingAdmin.java
index e7f8f439d51..71fc8102541 100644
--- a/clients/src/main/java/org/apache/kafka/clients/admin/ForwardingAdmin.java
+++ b/clients/src/main/java/org/apache/kafka/clients/admin/ForwardingAdmin.java
@@ -272,6 +272,11 @@ public class ForwardingAdmin implements Admin {
return delegate.unregisterBroker(brokerId, options);
}
+ @Override
+ public UnregisterControllerResult unregisterController(int controllerId,
UnregisterControllerOptions options) {
+ return delegate.unregisterController(controllerId, options);
+ }
+
@Override
public DescribeProducersResult
describeProducers(Collection<TopicPartition> partitions,
DescribeProducersOptions options) {
return delegate.describeProducers(partitions, options);
diff --git
a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
index 873c2a3e5c5..4f13eb232c3 100644
--- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
+++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
@@ -171,6 +171,7 @@ import org.apache.kafka.common.message.MetadataRequestData;
import org.apache.kafka.common.message.RemoveRaftVoterRequestData;
import org.apache.kafka.common.message.RenewDelegationTokenRequestData;
import org.apache.kafka.common.message.UnregisterBrokerRequestData;
+import org.apache.kafka.common.message.UnregisterControllerRequestData;
import org.apache.kafka.common.message.UpdateFeaturesRequestData;
import
org.apache.kafka.common.message.UpdateFeaturesResponseData.UpdatableFeatureResult;
import org.apache.kafka.common.metrics.KafkaMetric;
@@ -252,6 +253,8 @@ import
org.apache.kafka.common.requests.RenewDelegationTokenRequest;
import org.apache.kafka.common.requests.RenewDelegationTokenResponse;
import org.apache.kafka.common.requests.UnregisterBrokerRequest;
import org.apache.kafka.common.requests.UnregisterBrokerResponse;
+import org.apache.kafka.common.requests.UnregisterControllerRequest;
+import org.apache.kafka.common.requests.UnregisterControllerResponse;
import org.apache.kafka.common.requests.UpdateFeaturesRequest;
import org.apache.kafka.common.requests.UpdateFeaturesResponse;
import org.apache.kafka.common.security.auth.KafkaPrincipal;
@@ -4935,6 +4938,48 @@ public class KafkaAdminClient extends AdminClient {
return new UnregisterBrokerResult(future);
}
+ @Override
+ public UnregisterControllerResult unregisterController(int controllerId,
UnregisterControllerOptions options) {
+ final KafkaFutureImpl<Void> future = new KafkaFutureImpl<>();
+ final long now = time.milliseconds();
+ final Call call = new Call("unregisterController", calcDeadlineMs(now,
options.timeoutMs()),
+ new LeastLoadedBrokerOrActiveKController()) {
+
+ @Override
+ UnregisterControllerRequest.Builder createRequest(int timeoutMs) {
+ UnregisterControllerRequestData data =
+ new
UnregisterControllerRequestData().setControllerId(controllerId);
+ return new UnregisterControllerRequest.Builder(data);
+ }
+
+ @Override
+ void handleResponse(AbstractResponse abstractResponse) {
+ final UnregisterControllerResponse response =
+ (UnregisterControllerResponse) abstractResponse;
+ Errors error = Errors.forCode(response.data().errorCode());
+ switch (error) {
+ case NONE:
+ future.complete(null);
+ break;
+ case REQUEST_TIMED_OUT:
+ throw error.exception(response.data().errorMessage());
+ default:
+ log.error("Unregister controller request for
controller ID {} failed: {}",
+ controllerId, response.data().errorMessage());
+
future.completeExceptionally(error.exception(response.data().errorMessage()));
+ break;
+ }
+ }
+
+ @Override
+ void handleFailure(Throwable throwable) {
+ future.completeExceptionally(throwable);
+ }
+ };
+ runnable.call(call, now);
+ return new UnregisterControllerResult(future);
+ }
+
@Override
public DescribeProducersResult
describeProducers(Collection<TopicPartition> topicPartitions,
DescribeProducersOptions options) {
PartitionLeaderStrategy.PartitionLeaderFuture<DescribeProducersResult.PartitionProducerState>
future =
diff --git
a/clients/src/main/java/org/apache/kafka/clients/admin/UnregisterControllerOptions.java
b/clients/src/main/java/org/apache/kafka/clients/admin/UnregisterControllerOptions.java
new file mode 100644
index 00000000000..b626dd72125
--- /dev/null
+++
b/clients/src/main/java/org/apache/kafka/clients/admin/UnregisterControllerOptions.java
@@ -0,0 +1,27 @@
+/*
+ * 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.clients.admin;
+
+import org.apache.kafka.common.annotation.InterfaceAudience;
+
+/**
+ * Options for {@link Admin#unregisterController(int,
UnregisterControllerOptions)}.
+ */
[email protected]
+public class UnregisterControllerOptions extends
AbstractOptions<UnregisterControllerOptions> {
+}
diff --git
a/clients/src/main/java/org/apache/kafka/clients/admin/UnregisterControllerResult.java
b/clients/src/main/java/org/apache/kafka/clients/admin/UnregisterControllerResult.java
new file mode 100644
index 00000000000..07a907d4376
--- /dev/null
+++
b/clients/src/main/java/org/apache/kafka/clients/admin/UnregisterControllerResult.java
@@ -0,0 +1,42 @@
+/*
+ * 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.clients.admin;
+
+import org.apache.kafka.common.KafkaFuture;
+import org.apache.kafka.common.annotation.InterfaceAudience;
+
+/**
+ * The result of the {@link Admin#unregisterController(int,
UnregisterControllerOptions)} call.
+ *
+ * The API of this class is evolving, see {@link Admin} for details.
+ */
[email protected]
+public class UnregisterControllerResult {
+ private final KafkaFuture<Void> future;
+
+ UnregisterControllerResult(final KafkaFuture<Void> future) {
+ this.future = future;
+ }
+
+ /**
+ * Return a future which succeeds if the operation is successful.
+ */
+ public KafkaFuture<Void> all() {
+ return future;
+ }
+}
diff --git
a/clients/src/main/java/org/apache/kafka/common/errors/ControllerIdNotRegisteredException.java
b/clients/src/main/java/org/apache/kafka/common/errors/ControllerIdNotRegisteredException.java
new file mode 100644
index 00000000000..dbd6290c845
--- /dev/null
+++
b/clients/src/main/java/org/apache/kafka/common/errors/ControllerIdNotRegisteredException.java
@@ -0,0 +1,32 @@
+/*
+ * 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.errors;
+
+import org.apache.kafka.common.annotation.InterfaceAudience;
+
[email protected]
+public class ControllerIdNotRegisteredException extends ApiException {
+
+ public ControllerIdNotRegisteredException(String message) {
+ super(message);
+ }
+
+ public ControllerIdNotRegisteredException(String message, Throwable
throwable) {
+ super(message, throwable);
+ }
+
+}
diff --git
a/clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java
b/clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java
index 22dc37e67d1..5dbdd8f7cb3 100644
--- a/clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java
+++ b/clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java
@@ -138,7 +138,8 @@ public enum ApiKeys {
DESCRIBE_SHARE_GROUP_OFFSETS(ApiMessageType.DESCRIBE_SHARE_GROUP_OFFSETS),
ALTER_SHARE_GROUP_OFFSETS(ApiMessageType.ALTER_SHARE_GROUP_OFFSETS),
DELETE_SHARE_GROUP_OFFSETS(ApiMessageType.DELETE_SHARE_GROUP_OFFSETS),
-
STREAMS_GROUP_TOPOLOGY_DESCRIPTION_UPDATE(ApiMessageType.STREAMS_GROUP_TOPOLOGY_DESCRIPTION_UPDATE);
+
STREAMS_GROUP_TOPOLOGY_DESCRIPTION_UPDATE(ApiMessageType.STREAMS_GROUP_TOPOLOGY_DESCRIPTION_UPDATE),
+ UNREGISTER_CONTROLLER(ApiMessageType.UNREGISTER_CONTROLLER, false, true);
private static final Map<ApiMessageType.ListenerType, EnumSet<ApiKeys>>
APIS_BY_LISTENER =
new EnumMap<>(ApiMessageType.ListenerType.class);
diff --git a/clients/src/main/java/org/apache/kafka/common/protocol/Errors.java
b/clients/src/main/java/org/apache/kafka/common/protocol/Errors.java
index 49b67b738d5..5b49ed383cd 100644
--- a/clients/src/main/java/org/apache/kafka/common/protocol/Errors.java
+++ b/clients/src/main/java/org/apache/kafka/common/protocol/Errors.java
@@ -22,6 +22,7 @@ import
org.apache.kafka.common.errors.BrokerIdNotRegisteredException;
import org.apache.kafka.common.errors.BrokerNotAvailableException;
import org.apache.kafka.common.errors.ClusterAuthorizationException;
import org.apache.kafka.common.errors.ConcurrentTransactionsException;
+import org.apache.kafka.common.errors.ControllerIdNotRegisteredException;
import org.apache.kafka.common.errors.ControllerMovedException;
import org.apache.kafka.common.errors.CoordinatorLoadInProgressException;
import org.apache.kafka.common.errors.CoordinatorNotAvailableException;
@@ -422,7 +423,8 @@ public enum Errors {
STREAMS_TOPOLOGY_FENCED(132, "The supplied topology epoch is outdated.",
StreamsTopologyFencedException::new),
SHARE_SESSION_LIMIT_REACHED(133, "The limit of share sessions has been
reached.", ShareSessionLimitReachedException::new),
GROUP_DELETION_FAILED(134, "DeleteGroups could not complete; see the error
message on the per-group result for details.",
GroupDeletionFailedException::new),
- STREAMS_TOPOLOGY_DESCRIPTION_UPDATE_FAILED(135, "The broker could not
process the topology description update; see the error message for details.",
StreamsTopologyDescriptionUpdateFailedException::new);
+ STREAMS_TOPOLOGY_DESCRIPTION_UPDATE_FAILED(135, "The broker could not
process the topology description update; see the error message for details.",
StreamsTopologyDescriptionUpdateFailedException::new),
+ CONTROLLER_ID_NOT_REGISTERED(136, "The given controller ID was not
registered.", ControllerIdNotRegisteredException::new);
private static final Logger log = LoggerFactory.getLogger(Errors.class);
diff --git
a/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequest.java
b/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequest.java
index 8630c379ed6..58f735e12ad 100644
---
a/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequest.java
+++
b/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequest.java
@@ -356,6 +356,8 @@ public abstract class AbstractRequest implements
AbstractRequestResponse {
return DeleteShareGroupOffsetsRequest.parse(readable,
apiVersion);
case STREAMS_GROUP_TOPOLOGY_DESCRIPTION_UPDATE:
return
StreamsGroupTopologyDescriptionUpdateRequest.parse(readable, apiVersion);
+ case UNREGISTER_CONTROLLER:
+ return UnregisterControllerRequest.parse(readable, apiVersion);
default:
throw new AssertionError(String.format("ApiKey %s is not
currently handled in `parseRequest`, the " +
"code should be updated to do so.", apiKey));
diff --git
a/clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java
b/clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java
index bf265a58ca9..8d808a464b4 100644
---
a/clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java
+++
b/clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java
@@ -292,6 +292,8 @@ public abstract class AbstractResponse implements
AbstractRequestResponse {
return DeleteShareGroupOffsetsResponse.parse(readable,
version);
case STREAMS_GROUP_TOPOLOGY_DESCRIPTION_UPDATE:
return
StreamsGroupTopologyDescriptionUpdateResponse.parse(readable, version);
+ case UNREGISTER_CONTROLLER:
+ return UnregisterControllerResponse.parse(readable, version);
default:
throw new AssertionError(String.format("ApiKey %s is not
currently handled in `parseResponse`, the " +
"code should be updated to do so.", apiKey));
diff --git
a/clients/src/main/java/org/apache/kafka/common/requests/UnregisterControllerRequest.java
b/clients/src/main/java/org/apache/kafka/common/requests/UnregisterControllerRequest.java
new file mode 100644
index 00000000000..20d95bb0e42
--- /dev/null
+++
b/clients/src/main/java/org/apache/kafka/common/requests/UnregisterControllerRequest.java
@@ -0,0 +1,66 @@
+/*
+ * 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.requests;
+
+import org.apache.kafka.common.message.UnregisterControllerRequestData;
+import org.apache.kafka.common.message.UnregisterControllerResponseData;
+import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.Errors;
+import org.apache.kafka.common.protocol.Readable;
+
+public class UnregisterControllerRequest extends AbstractRequest {
+
+ public static class Builder extends
AbstractRequest.Builder<UnregisterControllerRequest> {
+ private final UnregisterControllerRequestData data;
+
+ public Builder(UnregisterControllerRequestData data) {
+ super(ApiKeys.UNREGISTER_CONTROLLER);
+ this.data = data;
+ }
+
+ @Override
+ public UnregisterControllerRequest build(short version) {
+ return new UnregisterControllerRequest(data, version);
+ }
+ }
+
+ private final UnregisterControllerRequestData data;
+
+ public UnregisterControllerRequest(UnregisterControllerRequestData data,
short version) {
+ super(ApiKeys.UNREGISTER_CONTROLLER, version);
+ this.data = data;
+ }
+
+ @Override
+ public UnregisterControllerRequestData data() {
+ return data;
+ }
+
+ @Override
+ public UnregisterControllerResponse getErrorResponse(int throttleTimeMs,
Throwable e) {
+ Errors error = Errors.forException(e);
+ return new UnregisterControllerResponse(new
UnregisterControllerResponseData()
+ .setThrottleTimeMs(throttleTimeMs)
+ .setErrorCode(error.code())
+ .setErrorMessage(e.getMessage()));
+ }
+
+ public static UnregisterControllerRequest parse(Readable readable, short
version) {
+ return new UnregisterControllerRequest(new
UnregisterControllerRequestData(readable, version),
+ version);
+ }
+}
diff --git
a/clients/src/main/java/org/apache/kafka/common/requests/UnregisterControllerResponse.java
b/clients/src/main/java/org/apache/kafka/common/requests/UnregisterControllerResponse.java
new file mode 100644
index 00000000000..e253da37c57
--- /dev/null
+++
b/clients/src/main/java/org/apache/kafka/common/requests/UnregisterControllerResponse.java
@@ -0,0 +1,63 @@
+/*
+ * 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.requests;
+
+import org.apache.kafka.common.message.UnregisterControllerResponseData;
+import org.apache.kafka.common.protocol.ApiKeys;
+import org.apache.kafka.common.protocol.Errors;
+import org.apache.kafka.common.protocol.Readable;
+
+import java.util.EnumMap;
+import java.util.Map;
+
+public class UnregisterControllerResponse extends AbstractResponse {
+ private final UnregisterControllerResponseData data;
+
+ public UnregisterControllerResponse(UnregisterControllerResponseData data)
{
+ super(ApiKeys.UNREGISTER_CONTROLLER);
+ this.data = data;
+ }
+
+ @Override
+ public UnregisterControllerResponseData data() {
+ return data;
+ }
+
+ @Override
+ public int throttleTimeMs() {
+ return data.throttleTimeMs();
+ }
+
+ @Override
+ public void maybeSetThrottleTimeMs(int throttleTimeMs) {
+ data.setThrottleTimeMs(throttleTimeMs);
+ }
+
+ @Override
+ public Map<Errors, Integer> errorCounts() {
+ Map<Errors, Integer> errorCounts = new EnumMap<>(Errors.class);
+ if (data.errorCode() != 0) {
+ errorCounts.put(Errors.forCode(data.errorCode()), 1);
+ }
+ return errorCounts;
+ }
+
+ public static UnregisterControllerResponse parse(Readable readable, short
version) {
+ return new UnregisterControllerResponse(new
UnregisterControllerResponseData(readable, version));
+ }
+}
diff --git
a/clients/src/main/resources/common/message/UnregisterControllerRequest.json
b/clients/src/main/resources/common/message/UnregisterControllerRequest.json
new file mode 100644
index 00000000000..39fdca54ebb
--- /dev/null
+++ b/clients/src/main/resources/common/message/UnregisterControllerRequest.json
@@ -0,0 +1,27 @@
+// 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.
+
+{
+ "apiKey": 94,
+ "type": "request",
+ "listeners": ["broker", "controller"],
+ "name": "UnregisterControllerRequest",
+ "validVersions": "0",
+ "flexibleVersions": "0+",
+ "fields": [
+ { "name": "ControllerId", "type": "int32", "versions": "0+",
+ "about": "The controller ID to unregister." }
+ ]
+}
diff --git
a/clients/src/main/resources/common/message/UnregisterControllerResponse.json
b/clients/src/main/resources/common/message/UnregisterControllerResponse.json
new file mode 100644
index 00000000000..0dfcabc0031
--- /dev/null
+++
b/clients/src/main/resources/common/message/UnregisterControllerResponse.json
@@ -0,0 +1,30 @@
+// 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.
+
+{
+ "apiKey": 94,
+ "type": "response",
+ "name": "UnregisterControllerResponse",
+ "validVersions": "0",
+ "flexibleVersions": "0+",
+ "fields": [
+ { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+",
+ "about": "Duration in milliseconds for which the request was throttled
due to a quota violation, or zero if the request did not violate any quota." },
+ { "name": "ErrorCode", "type": "int16", "versions": "0+",
+ "about": "The error code, or 0 if there was no error." },
+ { "name": "ErrorMessage", "type": "string", "versions": "0+",
"nullableVersions": "0+",
+ "about": "The top-level error message, or `null` if there was no
top-level error." }
+ ]
+}
diff --git
a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java
b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java
index 0162d2faa49..e90cfe4b2a5 100644
---
a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java
+++
b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java
@@ -166,6 +166,7 @@ import
org.apache.kafka.common.message.RemoveRaftVoterResponseData;
import org.apache.kafka.common.message.ShareGroupDescribeResponseData;
import org.apache.kafka.common.message.StreamsGroupDescribeResponseData;
import org.apache.kafka.common.message.UnregisterBrokerResponseData;
+import org.apache.kafka.common.message.UnregisterControllerResponseData;
import org.apache.kafka.common.message.WriteTxnMarkersResponseData;
import org.apache.kafka.common.protocol.ApiKeys;
import org.apache.kafka.common.protocol.Errors;
@@ -247,6 +248,7 @@ import org.apache.kafka.common.requests.RequestTestUtils;
import org.apache.kafka.common.requests.ShareGroupDescribeResponse;
import org.apache.kafka.common.requests.StreamsGroupDescribeResponse;
import org.apache.kafka.common.requests.UnregisterBrokerResponse;
+import org.apache.kafka.common.requests.UnregisterControllerResponse;
import org.apache.kafka.common.requests.UpdateFeaturesRequest;
import org.apache.kafka.common.requests.UpdateFeaturesResponse;
import org.apache.kafka.common.requests.WriteTxnMarkersRequest;
@@ -302,6 +304,7 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Function;
import java.util.function.Predicate;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@@ -10048,97 +10051,147 @@ public class KafkaAdminClientTest {
}
}
+ private static final int UNREGISTER_NODE_ID = 1;
+
+ private static final Function<Errors, AbstractResponse>
BROKER_RESPONSE_FACTORY =
+ error -> new UnregisterBrokerResponse(new
UnregisterBrokerResponseData()
+ .setErrorCode(error.code())
+ .setErrorMessage(error.message()));
+
+ private static final Function<Errors, AbstractResponse>
CONTROLLER_RESPONSE_FACTORY =
+ error -> new UnregisterControllerResponse(new
UnregisterControllerResponseData()
+ .setErrorCode(error.code())
+ .setErrorMessage(error.message()));
+
+ private static final Function<Admin, KafkaFuture<Void>>
UNREGISTER_BROKER_CALL =
+ admin -> admin.unregisterBroker(UNREGISTER_NODE_ID).all();
+
+ private static final Function<Admin, KafkaFuture<Void>>
UNREGISTER_CONTROLLER_CALL =
+ admin -> admin.unregisterController(UNREGISTER_NODE_ID).all();
+
+ private void runUnregisterScenario(
+ AdminClientUnitTestEnv env,
+ ApiKeys apiKey,
+ Function<Errors, AbstractResponse> responseFactory,
+ Function<Admin, KafkaFuture<Void>> adminCall,
+ List<Errors> responsesToPrepare,
+ Class<? extends Throwable> expectedException
+ ) throws ExecutionException, InterruptedException {
+ env.kafkaClient().setNodeApiVersions(
+ NodeApiVersions.create(apiKey.id, (short) 0, (short) 0));
+ for (Errors error : responsesToPrepare) {
+ env.kafkaClient().prepareResponse(responseFactory.apply(error));
+ }
+ KafkaFuture<Void> future = adminCall.apply(env.adminClient());
+ assertNotNull(future);
+ if (expectedException == null) {
+ future.get();
+ } else {
+ TestUtils.assertFutureThrows(expectedException, future);
+ }
+ }
+
@Test
public void testUnregisterBrokerSuccess() throws InterruptedException,
ExecutionException {
- int nodeId = 1;
- try (final AdminClientUnitTestEnv env = mockClientEnv()) {
- env.kafkaClient().setNodeApiVersions(
- NodeApiVersions.create(ApiKeys.UNREGISTER_BROKER.id,
(short) 0, (short) 0));
-
env.kafkaClient().prepareResponse(prepareUnregisterBrokerResponse(Errors.NONE,
0));
- UnregisterBrokerResult result =
env.adminClient().unregisterBroker(nodeId);
- // Validate response
- assertNotNull(result.all());
- result.all().get();
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ runUnregisterScenario(env, ApiKeys.UNREGISTER_BROKER,
BROKER_RESPONSE_FACTORY,
+ UNREGISTER_BROKER_CALL, List.of(Errors.NONE), null);
}
}
@Test
- public void testUnregisterBrokerFailure() {
- int nodeId = 1;
- try (final AdminClientUnitTestEnv env = mockClientEnv()) {
- env.kafkaClient().setNodeApiVersions(
- NodeApiVersions.create(ApiKeys.UNREGISTER_BROKER.id,
(short) 0, (short) 0));
-
env.kafkaClient().prepareResponse(prepareUnregisterBrokerResponse(Errors.UNKNOWN_SERVER_ERROR,
0));
- UnregisterBrokerResult result =
env.adminClient().unregisterBroker(nodeId);
- // Validate response
- assertNotNull(result.all());
- TestUtils.assertFutureThrows(UnknownServerException.class,
result.all());
+ public void testUnregisterBrokerFailure() throws ExecutionException,
InterruptedException {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ runUnregisterScenario(env, ApiKeys.UNREGISTER_BROKER,
BROKER_RESPONSE_FACTORY,
+ UNREGISTER_BROKER_CALL, List.of(Errors.UNKNOWN_SERVER_ERROR),
UnknownServerException.class);
}
}
@Test
public void testUnregisterBrokerTimeoutAndSuccessRetry() throws
ExecutionException, InterruptedException {
- int nodeId = 1;
- try (final AdminClientUnitTestEnv env = mockClientEnv()) {
- env.kafkaClient().setNodeApiVersions(
- NodeApiVersions.create(ApiKeys.UNREGISTER_BROKER.id,
(short) 0, (short) 0));
-
env.kafkaClient().prepareResponse(prepareUnregisterBrokerResponse(Errors.REQUEST_TIMED_OUT,
0));
-
env.kafkaClient().prepareResponse(prepareUnregisterBrokerResponse(Errors.NONE,
0));
-
- UnregisterBrokerResult result =
env.adminClient().unregisterBroker(nodeId);
-
- // Validate response
- assertNotNull(result.all());
- result.all().get();
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ runUnregisterScenario(env, ApiKeys.UNREGISTER_BROKER,
BROKER_RESPONSE_FACTORY,
+ UNREGISTER_BROKER_CALL, List.of(Errors.REQUEST_TIMED_OUT,
Errors.NONE), null);
}
}
@Test
- public void testUnregisterBrokerTimeoutAndFailureRetry() {
- int nodeId = 1;
- try (final AdminClientUnitTestEnv env = mockClientEnv()) {
- env.kafkaClient().setNodeApiVersions(
- NodeApiVersions.create(ApiKeys.UNREGISTER_BROKER.id,
(short) 0, (short) 0));
-
env.kafkaClient().prepareResponse(prepareUnregisterBrokerResponse(Errors.REQUEST_TIMED_OUT,
0));
-
env.kafkaClient().prepareResponse(prepareUnregisterBrokerResponse(Errors.UNKNOWN_SERVER_ERROR,
0));
+ public void testUnregisterBrokerTimeoutAndFailureRetry() throws
ExecutionException, InterruptedException {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ runUnregisterScenario(env, ApiKeys.UNREGISTER_BROKER,
BROKER_RESPONSE_FACTORY,
+ UNREGISTER_BROKER_CALL, List.of(Errors.REQUEST_TIMED_OUT,
Errors.UNKNOWN_SERVER_ERROR),
+ UnknownServerException.class);
+ }
+ }
- UnregisterBrokerResult result =
env.adminClient().unregisterBroker(nodeId);
+ @Test
+ public void testUnregisterBrokerTimeoutMaxRetry() throws
ExecutionException, InterruptedException {
+ try (AdminClientUnitTestEnv env = mockClientEnv(Time.SYSTEM,
AdminClientConfig.RETRIES_CONFIG, "1")) {
+ runUnregisterScenario(env, ApiKeys.UNREGISTER_BROKER,
BROKER_RESPONSE_FACTORY,
+ UNREGISTER_BROKER_CALL, List.of(Errors.REQUEST_TIMED_OUT,
Errors.REQUEST_TIMED_OUT),
+ TimeoutException.class);
+ }
+ }
- // Validate response
- assertNotNull(result.all());
- TestUtils.assertFutureThrows(UnknownServerException.class,
result.all());
+ @Test
+ public void testUnregisterBrokerTimeoutMaxWait() throws
ExecutionException, InterruptedException {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ runUnregisterScenario(env, ApiKeys.UNREGISTER_BROKER,
BROKER_RESPONSE_FACTORY,
+ admin -> admin.unregisterBroker(UNREGISTER_NODE_ID,
+ new UnregisterBrokerOptions().timeoutMs(10)).all(),
+ List.of(), TimeoutException.class);
}
}
@Test
- public void testUnregisterBrokerTimeoutMaxRetry() {
- int nodeId = 1;
- try (final AdminClientUnitTestEnv env = mockClientEnv(Time.SYSTEM,
AdminClientConfig.RETRIES_CONFIG, "1")) {
- env.kafkaClient().setNodeApiVersions(
- NodeApiVersions.create(ApiKeys.UNREGISTER_BROKER.id,
(short) 0, (short) 0));
-
env.kafkaClient().prepareResponse(prepareUnregisterBrokerResponse(Errors.REQUEST_TIMED_OUT,
0));
-
env.kafkaClient().prepareResponse(prepareUnregisterBrokerResponse(Errors.REQUEST_TIMED_OUT,
0));
+ public void testUnregisterControllerSuccess() throws InterruptedException,
ExecutionException {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ runUnregisterScenario(env, ApiKeys.UNREGISTER_CONTROLLER,
CONTROLLER_RESPONSE_FACTORY,
+ UNREGISTER_CONTROLLER_CALL, List.of(Errors.NONE), null);
+ }
+ }
- UnregisterBrokerResult result =
env.adminClient().unregisterBroker(nodeId);
+ @Test
+ public void testUnregisterControllerFailure() throws ExecutionException,
InterruptedException {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ runUnregisterScenario(env, ApiKeys.UNREGISTER_CONTROLLER,
CONTROLLER_RESPONSE_FACTORY,
+ UNREGISTER_CONTROLLER_CALL,
List.of(Errors.UNKNOWN_SERVER_ERROR), UnknownServerException.class);
+ }
+ }
- // Validate response
- assertNotNull(result.all());
- TestUtils.assertFutureThrows(TimeoutException.class, result.all());
+ @Test
+ public void testUnregisterControllerTimeoutAndSuccessRetry() throws
ExecutionException, InterruptedException {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ runUnregisterScenario(env, ApiKeys.UNREGISTER_CONTROLLER,
CONTROLLER_RESPONSE_FACTORY,
+ UNREGISTER_CONTROLLER_CALL, List.of(Errors.REQUEST_TIMED_OUT,
Errors.NONE), null);
}
}
@Test
- public void testUnregisterBrokerTimeoutMaxWait() {
- int nodeId = 1;
- try (final AdminClientUnitTestEnv env = mockClientEnv()) {
- env.kafkaClient().setNodeApiVersions(
- NodeApiVersions.create(ApiKeys.UNREGISTER_BROKER.id,
(short) 0, (short) 0));
+ public void testUnregisterControllerTimeoutAndFailureRetry() throws
ExecutionException, InterruptedException {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ runUnregisterScenario(env, ApiKeys.UNREGISTER_CONTROLLER,
CONTROLLER_RESPONSE_FACTORY,
+ UNREGISTER_CONTROLLER_CALL, List.of(Errors.REQUEST_TIMED_OUT,
Errors.UNKNOWN_SERVER_ERROR),
+ UnknownServerException.class);
+ }
+ }
- UnregisterBrokerResult result =
env.adminClient().unregisterBroker(nodeId, new
UnregisterBrokerOptions().timeoutMs(10));
+ @Test
+ public void testUnregisterControllerTimeoutMaxRetry() throws
ExecutionException, InterruptedException {
+ try (AdminClientUnitTestEnv env = mockClientEnv(Time.SYSTEM,
AdminClientConfig.RETRIES_CONFIG, "1")) {
+ runUnregisterScenario(env, ApiKeys.UNREGISTER_CONTROLLER,
CONTROLLER_RESPONSE_FACTORY,
+ UNREGISTER_CONTROLLER_CALL, List.of(Errors.REQUEST_TIMED_OUT,
Errors.REQUEST_TIMED_OUT),
+ TimeoutException.class);
+ }
+ }
- // Validate response
- assertNotNull(result.all());
- TestUtils.assertFutureThrows(TimeoutException.class, result.all());
+ @Test
+ public void testUnregisterControllerTimeoutMaxWait() throws
ExecutionException, InterruptedException {
+ try (AdminClientUnitTestEnv env = mockClientEnv()) {
+ runUnregisterScenario(env, ApiKeys.UNREGISTER_CONTROLLER,
CONTROLLER_RESPONSE_FACTORY,
+ admin -> admin.unregisterController(UNREGISTER_NODE_ID,
+ new UnregisterControllerOptions().timeoutMs(10)).all(),
+ List.of(), TimeoutException.class);
}
}
@@ -10824,13 +10877,6 @@ public class KafkaAdminClientTest {
admin.close();
}
- private UnregisterBrokerResponse prepareUnregisterBrokerResponse(Errors
error, int throttleTimeMs) {
- return new UnregisterBrokerResponse(new UnregisterBrokerResponseData()
- .setErrorCode(error.code())
- .setErrorMessage(error.message())
- .setThrottleTimeMs(throttleTimeMs));
- }
-
private DescribeLogDirsResponse prepareDescribeLogDirsResponse(Errors
error, String logDir) {
return new DescribeLogDirsResponse(new DescribeLogDirsResponseData()
.setResults(Collections.singletonList(
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 a8c0d5c88db..5b1f4e76c2d 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
@@ -242,6 +242,8 @@ import
org.apache.kafka.common.message.SyncGroupResponseData;
import org.apache.kafka.common.message.TxnOffsetCommitRequestData;
import org.apache.kafka.common.message.UnregisterBrokerRequestData;
import org.apache.kafka.common.message.UnregisterBrokerResponseData;
+import org.apache.kafka.common.message.UnregisterControllerRequestData;
+import org.apache.kafka.common.message.UnregisterControllerResponseData;
import org.apache.kafka.common.message.UpdateFeaturesRequestData;
import org.apache.kafka.common.message.UpdateFeaturesResponseData;
import org.apache.kafka.common.message.UpdateRaftVoterRequestData;
@@ -319,6 +321,7 @@ import static
org.apache.kafka.common.protocol.ApiKeys.PRODUCE;
import static org.apache.kafka.common.protocol.ApiKeys.SASL_AUTHENTICATE;
import static org.apache.kafka.common.protocol.ApiKeys.SYNC_GROUP;
import static org.apache.kafka.common.protocol.ApiKeys.UNREGISTER_BROKER;
+import static org.apache.kafka.common.protocol.ApiKeys.UNREGISTER_CONTROLLER;
import static
org.apache.kafka.common.requests.EndTxnRequest.LAST_STABLE_VERSION_BEFORE_TRANSACTION_V2;
import static
org.apache.kafka.common.requests.FetchMetadata.INVALID_SESSION_ID;
import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -347,6 +350,8 @@ public class RequestResponseTest {
toSkip.put(ELECT_LEADERS, List.of((short) 0));
// UnregisterBroker v0 contains the error message in the response
toSkip.put(UNREGISTER_BROKER, List.of((short) 0));
+ // UnregisterController v0 contains the error message in the response
+ toSkip.put(UNREGISTER_CONTROLLER, List.of((short) 0));
for (ApiKeys apikey : ApiKeys.values()) {
for (short version : apikey.allVersions()) {
@@ -929,16 +934,33 @@ public class RequestResponseTest {
UnregisterBrokerRequest request = new UnregisterBrokerRequest.Builder(
new UnregisterBrokerRequestData()
).build((short) 0);
- String customerErrorMessage = "customer error message";
+ String customErrorMessage = "custom error message";
UnregisterBrokerResponse response = request.getErrorResponse(
0,
- new RuntimeException(customerErrorMessage)
+ new RuntimeException(customErrorMessage)
);
assertEquals(0, response.throttleTimeMs());
assertEquals(Errors.UNKNOWN_SERVER_ERROR.code(),
response.data().errorCode());
- assertEquals(customerErrorMessage, response.data().errorMessage());
+ assertEquals(customErrorMessage, response.data().errorMessage());
+ }
+
+ @Test
+ public void testUnregisterControllerResponseWithUnknownServerError() {
+ UnregisterControllerRequest request = new
UnregisterControllerRequest.Builder(
+ new UnregisterControllerRequestData()
+ ).build((short) 0);
+ String customErrorMessage = "custom error message";
+
+ UnregisterControllerResponse response = request.getErrorResponse(
+ 0,
+ new RuntimeException(customErrorMessage)
+ );
+
+ assertEquals(0, response.throttleTimeMs());
+ assertEquals(Errors.UNKNOWN_SERVER_ERROR.code(),
response.data().errorCode());
+ assertEquals(customErrorMessage, response.data().errorMessage());
}
private ApiVersionsResponse defaultApiVersionsResponse() {
@@ -1145,6 +1167,7 @@ public class RequestResponseTest {
case ALTER_SHARE_GROUP_OFFSETS: return
createAlterShareGroupOffsetsRequest(version);
case DELETE_SHARE_GROUP_OFFSETS: return
createDeleteShareGroupOffsetsRequest(version);
case STREAMS_GROUP_TOPOLOGY_DESCRIPTION_UPDATE: return
createStreamsGroupTopologyDescriptionUpdateRequest(version);
+ case UNREGISTER_CONTROLLER: return
createUnregisterControllerRequest(version);
default: throw new IllegalArgumentException("Unknown API key " +
apikey);
}
}
@@ -1241,6 +1264,7 @@ public class RequestResponseTest {
case ALTER_SHARE_GROUP_OFFSETS: return
createAlterShareGroupOffsetsResponse();
case DELETE_SHARE_GROUP_OFFSETS: return
createDeleteShareGroupOffsetsResponse();
case STREAMS_GROUP_TOPOLOGY_DESCRIPTION_UPDATE: return
createStreamsGroupTopologyDescriptionUpdateResponse();
+ case UNREGISTER_CONTROLLER: return
createUnregisterControllerResponse();
default: throw new IllegalArgumentException("Unknown API key " +
apikey);
}
}
@@ -3630,6 +3654,15 @@ public class RequestResponseTest {
return new UnregisterBrokerResponse(new
UnregisterBrokerResponseData());
}
+ private UnregisterControllerRequest
createUnregisterControllerRequest(short version) {
+ UnregisterControllerRequestData data = new
UnregisterControllerRequestData().setControllerId(1);
+ return new UnregisterControllerRequest.Builder(data).build(version);
+ }
+
+ private UnregisterControllerResponse createUnregisterControllerResponse() {
+ return new UnregisterControllerResponse(new
UnregisterControllerResponseData());
+ }
+
private DescribeTransactionsRequest
createDescribeTransactionsRequest(short version) {
DescribeTransactionsRequestData data = new
DescribeTransactionsRequestData()
.setTransactionalIds(asList("t1", "t2", "t3"));
diff --git
a/clients/src/testFixtures/java/org/apache/kafka/clients/admin/MockAdminClient.java
b/clients/src/testFixtures/java/org/apache/kafka/clients/admin/MockAdminClient.java
index cc2a37b7aa9..70b2c1fde47 100644
---
a/clients/src/testFixtures/java/org/apache/kafka/clients/admin/MockAdminClient.java
+++
b/clients/src/testFixtures/java/org/apache/kafka/clients/admin/MockAdminClient.java
@@ -1368,6 +1368,17 @@ public class MockAdminClient extends AdminClient {
}
}
+ @Override
+ public UnregisterControllerResult unregisterController(int controllerId,
UnregisterControllerOptions options) {
+ if (usingRaftController) {
+ return new
UnregisterControllerResult(KafkaFuture.completedFuture(null));
+ } else {
+ KafkaFutureImpl<Void> future = new KafkaFutureImpl<>();
+ future.completeExceptionally(new UnsupportedVersionException(""));
+ return new UnregisterControllerResult(future);
+ }
+ }
+
@Override
public DescribeProducersResult
describeProducers(Collection<TopicPartition> partitions,
DescribeProducersOptions options) {
throw new UnsupportedOperationException("Not implemented yet");
diff --git a/core/src/main/scala/kafka/server/ControllerApis.scala
b/core/src/main/scala/kafka/server/ControllerApis.scala
index fc3de3e317d..909b113c578 100644
--- a/core/src/main/scala/kafka/server/ControllerApis.scala
+++ b/core/src/main/scala/kafka/server/ControllerApis.scala
@@ -108,6 +108,7 @@ class ControllerApis(
case ApiKeys.BROKER_REGISTRATION => handleBrokerRegistration(request)
case ApiKeys.BROKER_HEARTBEAT => handleBrokerHeartBeatRequest(request)
case ApiKeys.UNREGISTER_BROKER => handleUnregisterBroker(request)
+ case ApiKeys.UNREGISTER_CONTROLLER =>
handleUnregisterController(request)
case ApiKeys.ALTER_CLIENT_QUOTAS => handleAlterClientQuotas(request)
case ApiKeys.INCREMENTAL_ALTER_CONFIGS =>
handleIncrementalAlterConfigs(request)
case ApiKeys.ALTER_PARTITION_REASSIGNMENTS =>
handleAlterPartitionReassignments(request)
@@ -656,6 +657,27 @@ class ControllerApis(
}
}
+ def handleUnregisterController(request: Request): CompletableFuture[Unit] = {
+ val unregisterRequest = request.body(classOf[UnregisterControllerRequest])
+ authHelper.authorizeClusterOperation(request, ALTER)
+ val context = new ControllerRequestContext(request.context.header.data,
request.context.principal,
+ OptionalLong.empty())
+
+ controller.unregisterController(context,
unregisterRequest.data.controllerId).handle[Unit] { (_, e) =>
+ def createResponseCallback(requestThrottleMs: Int,
+ e: Throwable): UnregisterControllerResponse =
{
+ if (e != null) {
+ unregisterRequest.getErrorResponse(requestThrottleMs, e)
+ } else {
+ new UnregisterControllerResponse(new
UnregisterControllerResponseData().
+ setThrottleTimeMs(requestThrottleMs))
+ }
+ }
+ requestHelper.sendResponseMaybeThrottle(request,
+ requestThrottleMs => createResponseCallback(requestThrottleMs, e))
+ }
+ }
+
private def handleBrokerRegistration(request: Request):
CompletableFuture[Unit] = {
val registrationRequest = request.body(classOf[BrokerRegistrationRequest])
authHelper.authorizeClusterOperation(request, CLUSTER_ACTION)
diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala
b/core/src/main/scala/kafka/server/KafkaApis.scala
index d9d5879dac5..93af093142c 100644
--- a/core/src/main/scala/kafka/server/KafkaApis.scala
+++ b/core/src/main/scala/kafka/server/KafkaApis.scala
@@ -226,6 +226,7 @@ class KafkaApis(val requestChannel: RequestChannel,
case ApiKeys.DESCRIBE_CLUSTER => handleDescribeCluster(request)
case ApiKeys.DESCRIBE_PRODUCERS =>
handleDescribeProducersRequest(request)
case ApiKeys.UNREGISTER_BROKER => forwardToController(request)
+ case ApiKeys.UNREGISTER_CONTROLLER => forwardToController(request)
case ApiKeys.DESCRIBE_TRANSACTIONS =>
handleDescribeTransactionsRequest(request)
case ApiKeys.LIST_TRANSACTIONS =>
handleListTransactionsRequest(request)
case ApiKeys.DESCRIBE_QUORUM => forwardToController(request)
diff --git a/core/src/test/scala/unit/kafka/server/ControllerApisTest.scala
b/core/src/test/scala/unit/kafka/server/ControllerApisTest.scala
index 6a11bbfbe8e..2872d2fa227 100644
--- a/core/src/test/scala/unit/kafka/server/ControllerApisTest.scala
+++ b/core/src/test/scala/unit/kafka/server/ControllerApisTest.scala
@@ -435,6 +435,15 @@ class ControllerApisTest {
})
}
+ @Test
+ def testUnauthorizedHandleUnregisterController(): Unit = {
+ assertThrows(classOf[ClusterAuthorizationException], () => {
+ controllerApis = createControllerApis(Some(createDenyAllAuthorizer()),
new MockController.Builder().build())
+ controllerApis.handleUnregisterController(buildRequest(new
UnregisterControllerRequest.Builder(
+ new UnregisterControllerRequestData()).build(0)))
+ })
+ }
+
@Test
def testClose(): Unit = {
controllerApis = createControllerApis(Some(createDenyAllAuthorizer()),
mock(classOf[Controller]))
diff --git
a/core/src/test/scala/unit/kafka/server/ControllerRegistrationManagerTest.scala
b/core/src/test/scala/unit/kafka/server/ControllerRegistrationManagerTest.scala
index 833df332b1f..28c14b9e61d 100644
---
a/core/src/test/scala/unit/kafka/server/ControllerRegistrationManagerTest.scala
+++
b/core/src/test/scala/unit/kafka/server/ControllerRegistrationManagerTest.scala
@@ -19,7 +19,7 @@ package kafka.server
import org.apache.kafka.common.{Node, Uuid}
import org.apache.kafka.common.message.ControllerRegistrationResponseData
-import org.apache.kafka.common.metadata.{FeatureLevelRecord,
RegisterControllerRecord}
+import org.apache.kafka.common.metadata.{FeatureLevelRecord,
RegisterControllerRecord, UnregisterControllerRecord}
import org.apache.kafka.common.protocol.Errors
import org.apache.kafka.common.requests.ControllerRegistrationResponse
import org.apache.kafka.common.utils.internals.ExponentialBackoff
@@ -99,11 +99,16 @@ class ControllerRegistrationManagerTest {
failedAttempts.get(30, TimeUnit.SECONDS)
}
+ // This method simulates a metadata update by applying the given registration
+ // and unregistration modifiers to the previous image (e.g. an unregistration
+ // modifier of `_ => None` means do not unregister any nodes), and then
calling
+ // manager.onMetadataUpdate() with the resulting delta and new image.
private def doMetadataUpdate(
prevImage: MetadataImage,
manager: ControllerRegistrationManager,
metadataVersion: MetadataVersion,
- registrationModifier: RegisterControllerRecord =>
Option[RegisterControllerRecord]
+ registrationModifier: RegisterControllerRecord =>
Option[RegisterControllerRecord],
+ unregisterModifier: UnregisterControllerRecord =>
Option[UnregisterControllerRecord]
): MetadataImage = {
val delta = new MetadataDelta.Builder().
setImage(prevImage).
@@ -120,6 +125,13 @@ class ControllerRegistrationManagerTest {
}
}
}
+ if (metadataVersion.isControllerUnregistrationSupported) {
+ for (i <- Seq(1, 2, 3)) {
+
unregisterModifier(RecordTestUtils.createTestControllerUnregistration(i)).foreach
{
+ unregistration => delta.replay(unregistration)
+ }
+ }
+ }
val provenance = new MetadataProvenance(100, 200, 300, true)
val newImage = delta.apply(provenance)
val manifest = if
(!prevImage.features().metadataVersion().equals(Optional.of(metadataVersion))) {
@@ -181,7 +193,9 @@ class ControllerRegistrationManagerTest {
val image = doMetadataUpdate(MetadataImage.EMPTY,
manager,
metadataVersion,
- r => if (r.controllerId() == 1) None else Some(r))
+ r => if (r.controllerId() == 1) None else Some(r),
+ _ => None
+ )
if (!metadataVersionSupportsRegistration) {
assertFalse(registeredInLog(manager))
assertEquals((false, 0, 0), rpcStats(manager))
@@ -199,7 +213,9 @@ class ControllerRegistrationManagerTest {
doMetadataUpdate(image,
manager,
metadataVersion,
- r => Some(r))
+ r => Some(r),
+ _ => None
+ )
assertTrue(registeredInLog(manager))
}
} finally {
@@ -217,7 +233,9 @@ class ControllerRegistrationManagerTest {
doMetadataUpdate(MetadataImage.EMPTY,
manager,
MetadataVersion.IBP_3_7_IV0,
- r => Some(r.setIncarnationId(new Uuid(456, r.controllerId()))))
+ r => Some(r.setIncarnationId(new Uuid(456, r.controllerId()))),
+ _ => None
+ )
manager.start(context.mockChannelManager)
TestUtils.retryOnExceptionWithTimeout(30000, () => {
assertEquals((true, 0, 0), rpcStats(manager))
@@ -235,7 +253,9 @@ class ControllerRegistrationManagerTest {
doMetadataUpdate(MetadataImage.EMPTY,
manager,
MetadataVersion.IBP_3_7_IV0,
- r => Some(r.setIncarnationId(new Uuid(457, r.controllerId()))))
+ r => Some(r.setIncarnationId(new Uuid(457, r.controllerId()))),
+ _ => None
+ )
TestUtils.retryOnExceptionWithTimeout(30000, () => {
context.mockChannelManager.poll()
assertEquals((true, 1, 0), rpcStats(manager))
@@ -257,7 +277,9 @@ class ControllerRegistrationManagerTest {
doMetadataUpdate(MetadataImage.EMPTY,
manager,
MetadataVersion.IBP_3_7_IV0,
- r => if (r.controllerId() == 1) None else Some(r))
+ r => if (r.controllerId() == 1) None else Some(r),
+ _ => None
+ )
assertEquals((true, 0, 0), rpcStats(manager))
// the first response will trigger a retry
@@ -296,7 +318,9 @@ class ControllerRegistrationManagerTest {
doMetadataUpdate(MetadataImage.EMPTY,
manager,
MetadataVersion.IBP_3_7_IV0,
- r => if (r.controllerId() == 1) None else Some(r))
+ r => if (r.controllerId() == 1) None else Some(r),
+ _ => None
+ )
// pendingRpc = true, successfulRpcs = 0, failedRpcs = 0
assertEquals((true, 0, 0), rpcStats(manager))
assertEquals(1, context.mockChannelManager.unsentQueue.size())
@@ -318,4 +342,86 @@ class ControllerRegistrationManagerTest {
manager.close()
}
}
+
+ @Test
+ def testReRegistrationAfterUnregister(): Unit = {
+ val context = new RegistrationTestContext(configProperties)
+ val manager = newControllerRegistrationManager(context)
+ try {
+ context.controllerNodeProvider.node.set(controller1)
+ manager.start(context.mockChannelManager)
+
+ // initial state with controller registered in log
+ val image = doMetadataUpdate(MetadataImage.EMPTY,
+ manager,
+ MetadataVersion.IBP_4_4_IV2,
+ r => Some(r),
+ _ => None
+ )
+ assertTrue(registeredInLog(manager))
+
+ // Now unregister the controller via an UnregisterControllerRecord.
+ doMetadataUpdate(image,
+ manager,
+ MetadataVersion.IBP_4_4_IV2,
+ _ => None,
+ r => if (r.controllerId() == 1) Some(r) else None
+ )
+ assertFalse(registeredInLog(manager))
+
+ // The local manager should send a new registration RPC.
+ assertEquals((true, 0, 0), rpcStats(manager))
+ } finally {
+ manager.close()
+ }
+ }
+
+ @Test
+ def testReRegistrationAfterDifferentIncarnationId(): Unit = {
+ val context = new RegistrationTestContext(configProperties)
+ val manager = newControllerRegistrationManager(context)
+ try {
+ context.controllerNodeProvider.node.set(controller1)
+ manager.start(context.mockChannelManager)
+
+ // Register the controller.
+ val image = doMetadataUpdate(MetadataImage.EMPTY,
+ manager,
+ MetadataVersion.IBP_3_7_IV0,
+ r => if (r.controllerId() == 1) None else Some(r),
+ _ => None
+ )
+ assertEquals((true, 0, 0), rpcStats(manager))
+ context.mockClient.prepareResponseFrom(new
ControllerRegistrationResponse(
+ new ControllerRegistrationResponseData()), controller1)
+ context.mockChannelManager.poll()
+ assertEquals((false, 1, 0), rpcStats(manager))
+ assertFalse(registeredInLog(manager))
+
+ // Metadata update shows our registration is persisted.
+ val image2 = doMetadataUpdate(image,
+ manager,
+ MetadataVersion.IBP_3_7_IV0,
+ r => Some(r),
+ _ => None
+ )
+ assertTrue(registeredInLog(manager))
+
+ // Now a different incarnation registers with the same controller ID.
+ doMetadataUpdate(image2,
+ manager,
+ MetadataVersion.IBP_3_7_IV0,
+ r => if (r.controllerId() == 1) Some(r.setIncarnationId(new
Uuid(999999L, 1))) else Some(r),
+ _ => None
+ )
+
+ // The manager should detect the wrong incarnation and set
registeredInLog to false.
+ assertFalse(registeredInLog(manager))
+
+ // It should send a new registration RPC to reclaim its registration.
+ assertEquals((true, 1, 0), rpcStats(manager))
+ } finally {
+ manager.close()
+ }
+ }
}
diff --git a/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala
b/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala
index 02c523f0d3f..aca5832abc0 100644
--- a/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala
+++ b/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala
@@ -666,6 +666,9 @@ class RequestQuotaTest extends BaseRequestTest {
case ApiKeys.UNREGISTER_BROKER =>
new UnregisterBrokerRequest.Builder(new
UnregisterBrokerRequestData())
+ case ApiKeys.UNREGISTER_CONTROLLER =>
+ new UnregisterControllerRequest.Builder(new
UnregisterControllerRequestData())
+
case ApiKeys.DESCRIBE_TRANSACTIONS =>
new DescribeTransactionsRequest.Builder(new
DescribeTransactionsRequestData()
.setTransactionalIds(util.List.of("test-transactional-id")))
diff --git a/docs/getting-started/upgrade.md b/docs/getting-started/upgrade.md
index 61d32837bf2..d24bd01b728 100644
--- a/docs/getting-started/upgrade.md
+++ b/docs/getting-started/upgrade.md
@@ -53,6 +53,7 @@ type: docs
* Share groups now support dead-letter queue functionality as outlined in
[KIP-1191](https://cwiki.apache.org/confluence/x/fApJFg). Any records which are
released (beyond max delivery count) or rejected by the share consumer become
eligible for DLQ. Share group DLQ gets enabled when the Kafka feature
`share.version` is upgraded to 2. The user can configure a DLQ topic on a share
group by setting the dynamic config `errors.deadletterqueue.topic.name`
(default `""`) to the name of the DL [...]
* Kafka Connect distributed workers now support the
`internal.topics.automatic.creation.enable` configuration (default: `true`).
When set to `false`, Connect will not automatically create internal topics
(offset, config, status, and connector-specific offset topics) and will instead
fail at startup if any of these topics are missing. A new
`connect-internal-topics.sh` tool is also available for manually creating these
topics. For further details, please refer to [KIP-1209](https://cwik [...]
* Streams groups now support broker-side custom task assignors, registered
via the new broker configuration `group.streams.assignors` and selected per
group with the new group configuration `streams.assignor.name`. For further
details, please refer to
[KIP-1357](https://cwiki.apache.org/confluence/x/NoSnGQ).
+ * Controllers can now be unregistered from the cluster metadata. A new
`kafka-cluster.sh unregister-controller` command and a `--unregister` flag on
`kafka-metadata-quorum.sh remove-controller` are provided, backed by the new
`Admin#unregisterController` API and the new `UnregisterController` RPC. This
introduces the error code `CONTROLLER_ID_NOT_REGISTERED` (136) and requires
metadata version `4.4-IV2` (`IBP_4_4_IV2`). For further details, please refer
to [KIP-1312](https://cwiki.apac [...]
## Upgrading to 4.3.0
diff --git a/docs/operations/kraft.md b/docs/operations/kraft.md
index 3677a6a5e52..6af3245450b 100644
--- a/docs/operations/kraft.md
+++ b/docs/operations/kraft.md
@@ -232,6 +232,26 @@ When using controller endpoints use the
--bootstrap-controller flag:
$ bin/kafka-metadata-quorum.sh --bootstrap-controller localhost:9093
remove-controller --controller-id <id> --controller-directory-id <directory-id>
```
+To remove a KRaft voter and then unregister that controller from the cluster,
pass the `--unregister` flag:
+
+```bash
+$ bin/kafka-metadata-quorum.sh --bootstrap-server localhost:9092
remove-controller --controller-id <id> --controller-directory-id <directory-id>
--unregister
+```
+
+### Unregister Controller
+
+A controller that has already been removed from the voter set can be
unregistered from the cluster metadata using the `bin/kafka-cluster.sh
unregister-controller` command. The controller must be removed from the voter
set (and should be shut down) first — a running controller will re-register
itself. When using broker endpoints use the --bootstrap-server flag:
+
+```bash
+$ bin/kafka-cluster.sh unregister-controller --bootstrap-server localhost:9092
--id <id>
+```
+
+When using controller endpoints use the --bootstrap-controller flag:
+
+```bash
+$ bin/kafka-cluster.sh unregister-controller --bootstrap-controller
localhost:9093 --id <id>
+```
+
## Debugging
### Metadata Quorum Tool
diff --git a/docs/security/authorization-and-acls.md
b/docs/security/authorization-and-acls.md
index 0fc01c3ad48..965f9326ca9 100644
--- a/docs/security/authorization-and-acls.md
+++ b/docs/security/authorization-and-acls.md
@@ -2737,4 +2737,21 @@ Topic
<td>
+</td> </tr>
+<tr>
+<td>
+
+UNREGISTER_CONTROLLER (94)
+</td>
+<td>
+
+Alter
+</td>
+<td>
+
+Cluster
+</td>
+<td>
+
+
</td> </tr> </table>
diff --git
a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java
b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java
index b7665eeb8c0..686568c2c62 100644
---
a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java
+++
b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java
@@ -20,6 +20,7 @@ package org.apache.kafka.controller;
import org.apache.kafka.common.DirectoryId;
import org.apache.kafka.common.Uuid;
import org.apache.kafka.common.errors.BrokerIdNotRegisteredException;
+import org.apache.kafka.common.errors.ControllerIdNotRegisteredException;
import org.apache.kafka.common.errors.DuplicateBrokerRegistrationException;
import org.apache.kafka.common.errors.InconsistentClusterIdException;
import org.apache.kafka.common.errors.InvalidRegistrationException;
@@ -35,6 +36,7 @@ import
org.apache.kafka.common.metadata.RegisterControllerRecord;
import
org.apache.kafka.common.metadata.RegisterControllerRecord.ControllerFeatureCollection;
import org.apache.kafka.common.metadata.UnfenceBrokerRecord;
import org.apache.kafka.common.metadata.UnregisterBrokerRecord;
+import org.apache.kafka.common.metadata.UnregisterControllerRecord;
import org.apache.kafka.common.protocol.ApiMessage;
import org.apache.kafka.common.utils.Time;
import org.apache.kafka.common.utils.internals.LogContext;
@@ -498,6 +500,22 @@ public class ClusterControlManager {
return ControllerResult.atomicOf(records, null);
}
+ ControllerResult<Void> unregisterController(int controllerId) {
+ if
(!featureControl.metadataVersionOrThrow().isControllerUnregistrationSupported())
{
+ throw new UnsupportedVersionException("The current MetadataVersion
is too old to " +
+ "support controller unregistration.");
+ }
+ if (controllerRegistrations.get(controllerId) == null) {
+ throw new ControllerIdNotRegisteredException("Controller ID " +
controllerId +
+ " is not currently registered.");
+ }
+ List<ApiMessageAndVersion> records = new ArrayList<>();
+ records.add(new ApiMessageAndVersion(new UnregisterControllerRecord().
+ setControllerId(controllerId),
+ (short) 0));
+ return ControllerResult.atomicOf(records, null);
+ }
+
BrokerFeature processRegistrationFeature(
int brokerId,
FinalizedControllerFeatures finalizedFeatures,
@@ -610,6 +628,17 @@ public class ClusterControlManager {
}
}
+ public void replay(UnregisterControllerRecord record) {
+ int controllerId = record.controllerId();
+ ControllerRegistration registration =
controllerRegistrations.remove(controllerId);
+ if (registration == null) {
+ throw new RuntimeException(String.format("Unable to replay %s: no
controller " +
+ "registration found for that id", record));
+ } else {
+ log.info("Replayed {}", record);
+ }
+ }
+
public void replay(FenceBrokerRecord record) {
replayRegistrationChange(
record,
diff --git a/metadata/src/main/java/org/apache/kafka/controller/Controller.java
b/metadata/src/main/java/org/apache/kafka/controller/Controller.java
index 28aa076291b..c02e357ddf7 100644
--- a/metadata/src/main/java/org/apache/kafka/controller/Controller.java
+++ b/metadata/src/main/java/org/apache/kafka/controller/Controller.java
@@ -412,6 +412,20 @@ public interface Controller extends AclMutator,
AutoCloseable {
ControllerRegistrationRequestData request
);
+ /**
+ * Attempt to unregister the given controller.
+ *
+ * @param context The controller request context.
+ * @param controllerId The controller id to unregister.
+ *
+ * @return A future that is completed successfully when the
controller is
+ * unregistered.
+ */
+ CompletableFuture<Void> unregisterController(
+ ControllerRequestContext context,
+ int controllerId
+ );
+
/**
* Assign replicas to directories.
*
diff --git
a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java
b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java
index 6620c7e6dc9..128bab401fb 100644
--- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java
+++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java
@@ -82,6 +82,7 @@ import
org.apache.kafka.common.metadata.RemoveUserScramCredentialRecord;
import org.apache.kafka.common.metadata.TopicRecord;
import org.apache.kafka.common.metadata.UnfenceBrokerRecord;
import org.apache.kafka.common.metadata.UnregisterBrokerRecord;
+import org.apache.kafka.common.metadata.UnregisterControllerRecord;
import org.apache.kafka.common.metadata.UserScramCredentialRecord;
import org.apache.kafka.common.protocol.ApiMessage;
import org.apache.kafka.common.quota.ClientQuotaAlteration;
@@ -1325,6 +1326,9 @@ public final class QuorumController implements Controller
{
case REGISTER_CONTROLLER_RECORD:
clusterControl.replay((RegisterControllerRecord) message);
break;
+ case UNREGISTER_CONTROLLER_RECORD:
+ clusterControl.replay((UnregisterControllerRecord) message);
+ break;
case CLEAR_ELR_RECORD:
replicationControl.replay((ClearElrRecord) message);
break;
@@ -2153,6 +2157,21 @@ public final class QuorumController implements
Controller {
EnumSet.noneOf(ControllerOperationFlag.class));
}
+ @Override
+ public CompletableFuture<Void> unregisterController(
+ ControllerRequestContext context,
+ int controllerId
+ ) {
+ return appendWriteEvent("unregisterController", context.deadlineNs(),
+ () -> {
+ if (nodeId == controllerId) {
+ throw new InvalidRequestException("Controller cannot
unregister itself while it is active.");
+ }
+ return clusterControl.unregisterController(controllerId);
+ },
+ EnumSet.noneOf(ControllerOperationFlag.class));
+ }
+
@Override
public CompletableFuture<List<AclCreateResult>> createAcls(
ControllerRequestContext context,
diff --git a/metadata/src/main/java/org/apache/kafka/image/ClusterDelta.java
b/metadata/src/main/java/org/apache/kafka/image/ClusterDelta.java
index 1ab6e44423c..a8803853b6b 100644
--- a/metadata/src/main/java/org/apache/kafka/image/ClusterDelta.java
+++ b/metadata/src/main/java/org/apache/kafka/image/ClusterDelta.java
@@ -24,6 +24,7 @@ import org.apache.kafka.common.metadata.RegisterBrokerRecord;
import org.apache.kafka.common.metadata.RegisterControllerRecord;
import org.apache.kafka.common.metadata.UnfenceBrokerRecord;
import org.apache.kafka.common.metadata.UnregisterBrokerRecord;
+import org.apache.kafka.common.metadata.UnregisterControllerRecord;
import org.apache.kafka.metadata.BrokerRegistration;
import org.apache.kafka.metadata.BrokerRegistrationFencingChange;
import org.apache.kafka.metadata.BrokerRegistrationInControlledShutdownChange;
@@ -96,6 +97,10 @@ public final class ClusterDelta {
changedControllers.put(controller.id(), Optional.of(controller));
}
+ public void replay(UnregisterControllerRecord record) {
+ changedControllers.put(record.controllerId(), Optional.empty());
+ }
+
private BrokerRegistration getBrokerOrThrow(int brokerId, long epoch,
String action) {
BrokerRegistration broker = broker(brokerId);
if (broker == null) {
diff --git a/metadata/src/main/java/org/apache/kafka/image/MetadataDelta.java
b/metadata/src/main/java/org/apache/kafka/image/MetadataDelta.java
index 36d68dff5a2..50cc6c089dd 100644
--- a/metadata/src/main/java/org/apache/kafka/image/MetadataDelta.java
+++ b/metadata/src/main/java/org/apache/kafka/image/MetadataDelta.java
@@ -38,6 +38,7 @@ import
org.apache.kafka.common.metadata.RemoveUserScramCredentialRecord;
import org.apache.kafka.common.metadata.TopicRecord;
import org.apache.kafka.common.metadata.UnfenceBrokerRecord;
import org.apache.kafka.common.metadata.UnregisterBrokerRecord;
+import org.apache.kafka.common.metadata.UnregisterControllerRecord;
import org.apache.kafka.common.metadata.UserScramCredentialRecord;
import org.apache.kafka.common.protocol.ApiMessage;
import org.apache.kafka.metadata.SupportedConfigChecker;
@@ -267,6 +268,9 @@ public final class MetadataDelta {
case REGISTER_CONTROLLER_RECORD:
replay((RegisterControllerRecord) record);
break;
+ case UNREGISTER_CONTROLLER_RECORD:
+ replay((UnregisterControllerRecord) record);
+ break;
default:
throw new RuntimeException("Unknown metadata record type " +
type);
}
@@ -368,6 +372,10 @@ public final class MetadataDelta {
getOrCreateClusterDelta().replay(record);
}
+ public void replay(UnregisterControllerRecord record) {
+ getOrCreateClusterDelta().replay(record);
+ }
+
/**
* Create removal deltas for anything which was in the base image, but
which was not
* referenced in the snapshot records we just applied.
diff --git
a/metadata/src/main/resources/common/metadata/UnregisterControllerRecord.json
b/metadata/src/main/resources/common/metadata/UnregisterControllerRecord.json
new file mode 100644
index 00000000000..a6b6b44ec3d
--- /dev/null
+++
b/metadata/src/main/resources/common/metadata/UnregisterControllerRecord.json
@@ -0,0 +1,26 @@
+// 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.
+
+{
+ "apiKey": 29,
+ "type": "metadata",
+ "name": "UnregisterControllerRecord",
+ "validVersions": "0",
+ "flexibleVersions": "0+",
+ "fields": [
+ { "name": "ControllerId", "type": "int32", "versions": "0+",
+ "about": "The controller id." }
+ ]
+}
diff --git
a/metadata/src/test/java/org/apache/kafka/controller/ClusterControlManagerTest.java
b/metadata/src/test/java/org/apache/kafka/controller/ClusterControlManagerTest.java
index e9d4a0c8188..5f8b07527ba 100644
---
a/metadata/src/test/java/org/apache/kafka/controller/ClusterControlManagerTest.java
+++
b/metadata/src/test/java/org/apache/kafka/controller/ClusterControlManagerTest.java
@@ -20,6 +20,7 @@ package org.apache.kafka.controller;
import org.apache.kafka.common.DirectoryId;
import org.apache.kafka.common.Endpoint;
import org.apache.kafka.common.Uuid;
+import org.apache.kafka.common.errors.ControllerIdNotRegisteredException;
import org.apache.kafka.common.errors.DuplicateBrokerRegistrationException;
import org.apache.kafka.common.errors.InconsistentClusterIdException;
import org.apache.kafka.common.errors.InvalidRegistrationException;
@@ -35,6 +36,7 @@ import
org.apache.kafka.common.metadata.RegisterBrokerRecord.BrokerEndpoint;
import
org.apache.kafka.common.metadata.RegisterBrokerRecord.BrokerEndpointCollection;
import org.apache.kafka.common.metadata.UnfenceBrokerRecord;
import org.apache.kafka.common.metadata.UnregisterBrokerRecord;
+import org.apache.kafka.common.metadata.UnregisterControllerRecord;
import org.apache.kafka.common.security.auth.SecurityProtocol;
import org.apache.kafka.common.utils.MockTime;
import org.apache.kafka.common.utils.internals.LogContext;
@@ -1139,6 +1141,87 @@ public class ClusterControlManagerTest {
});
}
+ @Test
+ public void testUnregisterController() {
+ SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new
LogContext());
+ FeatureControlManager featureControl = new
FeatureControlManager.Builder().
+ setSnapshotRegistry(snapshotRegistry).
+ setQuorumFeatures(new QuorumFeatures(0,
+ QuorumFeatures.defaultSupportedFeatureMap(true),
+ List.of(0))).
+ build();
+ featureControl.replay(new FeatureLevelRecord().
+ setName(MetadataVersion.FEATURE_NAME).
+ setFeatureLevel(MetadataVersion.IBP_4_4_IV2.featureLevel()));
+ ClusterControlManager clusterControl = new
ClusterControlManager.Builder().
+ setTime(new MockTime(0, 0, 0)).
+ setSnapshotRegistry(snapshotRegistry).
+ setSessionTimeoutNs(1000).
+ setFeatureControlManager(featureControl).
+ setBrokerShutdownHandler((brokerId, isCleanShutdown, records) -> {
}).
+ build();
+ clusterControl.activate();
+
+ // Register a controller
+ ControllerResult<Void> registerResult =
clusterControl.registerController(
+ new ControllerRegistrationRequestData().setControllerId(1));
+ RecordTestUtils.replayAll(clusterControl, registerResult.records());
+ assertTrue(clusterControl.controllerRegistrations().containsKey(1));
+
+ // Unregister the controller
+ ControllerResult<Void> unregisterResult =
clusterControl.unregisterController(1);
+ assertEquals(1, unregisterResult.records().size());
+ RecordTestUtils.replayAll(clusterControl, unregisterResult.records());
+ assertFalse(clusterControl.controllerRegistrations().containsKey(1));
+ }
+
+ @Test
+ public void testUnregisterUnknownController() {
+ SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new
LogContext());
+ FeatureControlManager featureControl = new
FeatureControlManager.Builder().
+ setSnapshotRegistry(snapshotRegistry).
+ setQuorumFeatures(new QuorumFeatures(0,
+ QuorumFeatures.defaultSupportedFeatureMap(true),
+ List.of(0))).
+ build();
+ featureControl.replay(new FeatureLevelRecord().
+ setName(MetadataVersion.FEATURE_NAME).
+ setFeatureLevel(MetadataVersion.IBP_4_4_IV2.featureLevel()));
+ ClusterControlManager clusterControl = new
ClusterControlManager.Builder().
+ setTime(new MockTime(0, 0, 0)).
+ setSnapshotRegistry(snapshotRegistry).
+ setSessionTimeoutNs(1000).
+ setFeatureControlManager(featureControl).
+ setBrokerShutdownHandler((brokerId, isCleanShutdown, records) -> {
}).
+ build();
+ clusterControl.activate();
+
+ // Trying to unregister a non-registered controller should throw
ApiException
+ assertThrows(ControllerIdNotRegisteredException.class,
+ () -> clusterControl.unregisterController(1));
+
+ // Replaying unregister record for unknown controller should throw
RuntimeException
+ assertThrows(RuntimeException.class,
+ () -> clusterControl.replay(new
UnregisterControllerRecord().setControllerId(1)));
+ }
+
+ @Test
+ public void testUnregisterControllerWithUnsupportedMetadataVersion() {
+ FeatureControlManager featureControl = new
FeatureControlManager.Builder().
+ build();
+ featureControl.replay(new FeatureLevelRecord().
+ setName(MetadataVersion.FEATURE_NAME).
+ setFeatureLevel(MetadataVersion.MINIMUM_VERSION.featureLevel()));
+ ClusterControlManager clusterControl = new
ClusterControlManager.Builder().
+ setClusterId("fPZv1VBsRFmnlRvmGcOW9w").
+ setFeatureControlManager(featureControl).
+ setBrokerShutdownHandler((brokerId, isCleanShutdown, records)
-> { }).
+ build();
+ clusterControl.activate();
+ assertEquals("The current MetadataVersion is too old to support
controller unregistration.",
+ assertThrows(UnsupportedVersionException.class, () ->
clusterControl.unregisterController(1)).getMessage());
+ }
+
private FeatureControlManager createFeatureControlManager() {
FeatureControlManager featureControlManager = new
FeatureControlManager.Builder().build();
featureControlManager.replay(new FeatureLevelRecord().
diff --git
a/metadata/src/testFixtures/java/org/apache/kafka/metadata/RecordTestUtils.java
b/metadata/src/testFixtures/java/org/apache/kafka/metadata/RecordTestUtils.java
index c3f83050ad0..82e7e795e4c 100644
---
a/metadata/src/testFixtures/java/org/apache/kafka/metadata/RecordTestUtils.java
+++
b/metadata/src/testFixtures/java/org/apache/kafka/metadata/RecordTestUtils.java
@@ -20,6 +20,7 @@ package org.apache.kafka.metadata;
import org.apache.kafka.common.Uuid;
import org.apache.kafka.common.metadata.RegisterControllerRecord;
import org.apache.kafka.common.metadata.TopicRecord;
+import org.apache.kafka.common.metadata.UnregisterControllerRecord;
import org.apache.kafka.common.protocol.ApiMessage;
import org.apache.kafka.common.protocol.Message;
import org.apache.kafka.common.security.auth.SecurityProtocol;
@@ -275,4 +276,8 @@ public class RecordTestUtils {
)
));
}
+
+ public static UnregisterControllerRecord
createTestControllerUnregistration(int id) {
+ return new UnregisterControllerRecord().setControllerId(id);
+ }
}
diff --git
a/server-common/src/main/java/org/apache/kafka/server/common/MetadataVersion.java
b/server-common/src/main/java/org/apache/kafka/server/common/MetadataVersion.java
index b68407be06d..a6adf39b3e6 100644
---
a/server-common/src/main/java/org/apache/kafka/server/common/MetadataVersion.java
+++
b/server-common/src/main/java/org/apache/kafka/server/common/MetadataVersion.java
@@ -135,7 +135,10 @@ public enum MetadataVersion {
IBP_4_4_IV0(31, "4.4", "IV0", false),
// Add support for CIDR-based ACL host patterns (KIP-1276).
- IBP_4_4_IV1(32, "4.4", "IV1", true);
+ IBP_4_4_IV1(32, "4.4", "IV1", true),
+
+ // Add support for controller unregistration (KIP-1312).
+ IBP_4_4_IV2(33, "4.4", "IV2", true);
// NOTES when adding a new version:
// Update the default version in @ClusterTest annotation to point to the
latest version
@@ -254,6 +257,10 @@ public enum MetadataVersion {
return this.isAtLeast(MetadataVersion.IBP_3_7_IV0);
}
+ public boolean isControllerUnregistrationSupported() {
+ return this.isAtLeast(MetadataVersion.IBP_4_4_IV2);
+ }
+
public short partitionChangeRecordVersion() {
if (isElrSupported()) {
return (short) 2;
diff --git
a/server/src/main/java/org/apache/kafka/network/RequestConvertToJson.java
b/server/src/main/java/org/apache/kafka/network/RequestConvertToJson.java
index 023ef02a5f6..fe6b1b206e3 100644
--- a/server/src/main/java/org/apache/kafka/network/RequestConvertToJson.java
+++ b/server/src/main/java/org/apache/kafka/network/RequestConvertToJson.java
@@ -187,6 +187,8 @@ import
org.apache.kafka.common.message.TxnOffsetCommitRequestDataJsonConverter;
import
org.apache.kafka.common.message.TxnOffsetCommitResponseDataJsonConverter;
import
org.apache.kafka.common.message.UnregisterBrokerRequestDataJsonConverter;
import
org.apache.kafka.common.message.UnregisterBrokerResponseDataJsonConverter;
+import
org.apache.kafka.common.message.UnregisterControllerRequestDataJsonConverter;
+import
org.apache.kafka.common.message.UnregisterControllerResponseDataJsonConverter;
import org.apache.kafka.common.message.UpdateFeaturesRequestDataJsonConverter;
import org.apache.kafka.common.message.UpdateFeaturesResponseDataJsonConverter;
import org.apache.kafka.common.message.UpdateRaftVoterRequestDataJsonConverter;
@@ -372,6 +374,8 @@ import
org.apache.kafka.common.requests.TxnOffsetCommitRequest;
import org.apache.kafka.common.requests.TxnOffsetCommitResponse;
import org.apache.kafka.common.requests.UnregisterBrokerRequest;
import org.apache.kafka.common.requests.UnregisterBrokerResponse;
+import org.apache.kafka.common.requests.UnregisterControllerRequest;
+import org.apache.kafka.common.requests.UnregisterControllerResponse;
import org.apache.kafka.common.requests.UpdateFeaturesRequest;
import org.apache.kafka.common.requests.UpdateFeaturesResponse;
import org.apache.kafka.common.requests.UpdateRaftVoterRequest;
@@ -564,6 +568,8 @@ public class RequestConvertToJson {
TxnOffsetCommitRequestDataJsonConverter.write(((TxnOffsetCommitRequest)
request).data(), request.version());
case UNREGISTER_BROKER ->
UnregisterBrokerRequestDataJsonConverter.write(((UnregisterBrokerRequest)
request).data(), request.version());
+ case UNREGISTER_CONTROLLER ->
+
UnregisterControllerRequestDataJsonConverter.write(((UnregisterControllerRequest)
request).data(), request.version());
case UPDATE_FEATURES ->
UpdateFeaturesRequestDataJsonConverter.write(((UpdateFeaturesRequest)
request).data(), request.version());
case UPDATE_RAFT_VOTER ->
@@ -743,6 +749,8 @@ public class RequestConvertToJson {
TxnOffsetCommitResponseDataJsonConverter.write(((TxnOffsetCommitResponse)
response).data(), version);
case UNREGISTER_BROKER ->
UnregisterBrokerResponseDataJsonConverter.write(((UnregisterBrokerResponse)
response).data(), version);
+ case UNREGISTER_CONTROLLER ->
+
UnregisterControllerResponseDataJsonConverter.write(((UnregisterControllerResponse)
response).data(), version);
case UPDATE_FEATURES ->
UpdateFeaturesResponseDataJsonConverter.write(((UpdateFeaturesResponse)
response).data(), version);
case UPDATE_RAFT_VOTER ->
diff --git a/server/src/test/java/org/apache/kafka/server/KRaftClusterTest.java
b/server/src/test/java/org/apache/kafka/server/KRaftClusterTest.java
index cca9562d8cc..52ef977c634 100644
--- a/server/src/test/java/org/apache/kafka/server/KRaftClusterTest.java
+++ b/server/src/test/java/org/apache/kafka/server/KRaftClusterTest.java
@@ -40,11 +40,14 @@ import org.apache.kafka.common.Node;
import org.apache.kafka.common.Reconfigurable;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.TopicPartitionInfo;
+import org.apache.kafka.common.Uuid;
import org.apache.kafka.common.acl.AclBinding;
import org.apache.kafka.common.acl.AclBindingFilter;
import org.apache.kafka.common.config.ConfigResource;
import org.apache.kafka.common.config.ConfigResource.Type;
+import org.apache.kafka.common.errors.ControllerIdNotRegisteredException;
import org.apache.kafka.common.errors.InvalidPartitionsException;
+import org.apache.kafka.common.errors.InvalidRequestException;
import org.apache.kafka.common.errors.PolicyViolationException;
import org.apache.kafka.common.errors.UnsupportedVersionException;
import org.apache.kafka.common.message.DescribeClusterRequestData;
@@ -127,6 +130,7 @@ import java.util.stream.IntStream;
import java.util.stream.Stream;
import static org.apache.kafka.server.IntegrationTestUtils.connectAndReceive;
+import static org.apache.kafka.test.TestUtils.assertFutureThrows;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -675,6 +679,73 @@ public class KRaftClusterTest {
}
}
+ @ParameterizedTest
+ @ValueSource(booleans = {true, false})
+ public void testUnregisterController(boolean usingBootstrapControllers)
throws Exception {
+ final var nodes = new TestKitNodes.Builder().
+ setNumBrokerNodes(3).
+ setNumControllerNodes(3).
+ build();
+ final Map<Integer, Uuid> initialVoters = new HashMap<>();
+ for (final var controllerNode : nodes.controllerNodes().values()) {
+ initialVoters.put(
+ controllerNode.id(),
+ controllerNode.metadataDirectoryId()
+ );
+ }
+
+ try (KafkaClusterTestKit cluster = new
KafkaClusterTestKit.Builder(nodes).
+ setInitialVoterSet(initialVoters).
+ build()
+ ) {
+ cluster.format();
+ cluster.startup();
+ int controllerIdToUnregister =
cluster.controllers().keySet().iterator().next();
+ cluster.controllers().get(controllerIdToUnregister).shutdown();
+ cluster.waitForActiveController();
+
+ try (Admin admin = createAdminClient(cluster,
usingBootstrapControllers)) {
+ assertDoesNotThrow(() ->
admin.unregisterController(controllerIdToUnregister).all().get());
+ }
+
+ TestUtils.waitForCondition(() -> !clusterImage(cluster,
1).controllers().containsKey(controllerIdToUnregister),
+ "Timed out waiting for controller to be unregistered.");
+ }
+ }
+
+ @ParameterizedTest
+ @ValueSource(booleans = {true, false})
+ public void testUnregisterControllerError(boolean
usingBootstrapControllers) throws Exception {
+ try (KafkaClusterTestKit cluster = new KafkaClusterTestKit.Builder(
+ new TestKitNodes.Builder()
+ .setNumBrokerNodes(3)
+ .setNumControllerNodes(3)
+ .build()).build()) {
+ cluster.format();
+ cluster.startup();
+ int unknownId = 9999;
+ cluster.waitForActiveController();
+ int activeId = cluster.controllers().entrySet().stream()
+ .filter(e -> e.getValue().controller().isActive())
+ .map(Map.Entry::getKey)
+ .findFirst()
+ .orElseThrow();
+
+ try (Admin admin = createAdminClient(cluster,
usingBootstrapControllers)) {
+ assertFutureThrows(
+ ControllerIdNotRegisteredException.class,
+ admin.unregisterController(unknownId).all(),
+ String.format("Controller ID %s is not currently
registered.", unknownId)
+ );
+ assertFutureThrows(
+ InvalidRequestException.class,
+ admin.unregisterController(activeId).all(),
+ "Controller cannot unregister itself while it is active."
+ );
+ }
+ }
+ }
+
private ClusterImage clusterImage(KafkaClusterTestKit cluster, int
brokerId) {
return
cluster.brokers().get(brokerId).metadataCache().currentImage().cluster();
}
diff --git
a/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/TestingMetricsInterceptingAdminClient.java
b/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/TestingMetricsInterceptingAdminClient.java
index ae8e9f0924c..41698bd25bf 100644
---
a/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/TestingMetricsInterceptingAdminClient.java
+++
b/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/TestingMetricsInterceptingAdminClient.java
@@ -149,6 +149,8 @@ import
org.apache.kafka.clients.admin.TerminateTransactionOptions;
import org.apache.kafka.clients.admin.TerminateTransactionResult;
import org.apache.kafka.clients.admin.UnregisterBrokerOptions;
import org.apache.kafka.clients.admin.UnregisterBrokerResult;
+import org.apache.kafka.clients.admin.UnregisterControllerOptions;
+import org.apache.kafka.clients.admin.UnregisterControllerResult;
import org.apache.kafka.clients.admin.UpdateFeaturesOptions;
import org.apache.kafka.clients.admin.UpdateFeaturesResult;
import org.apache.kafka.clients.admin.UserScramCredentialAlteration;
@@ -408,6 +410,11 @@ public class TestingMetricsInterceptingAdminClient extends
AdminClient {
return adminDelegate.unregisterBroker(brokerId, options);
}
+ @Override
+ public UnregisterControllerResult unregisterController(final int
controllerId, final UnregisterControllerOptions options) {
+ return adminDelegate.unregisterController(controllerId, options);
+ }
+
@Override
public DescribeProducersResult describeProducers(final
Collection<TopicPartition> partitions, final DescribeProducersOptions options) {
return adminDelegate.describeProducers(partitions, options);
diff --git
a/test-common/test-common-internal-api/src/main/java/org/apache/kafka/common/test/api/ClusterTest.java
b/test-common/test-common-internal-api/src/main/java/org/apache/kafka/common/test/api/ClusterTest.java
index b884ec5734e..760050c172c 100644
---
a/test-common/test-common-internal-api/src/main/java/org/apache/kafka/common/test/api/ClusterTest.java
+++
b/test-common/test-common-internal-api/src/main/java/org/apache/kafka/common/test/api/ClusterTest.java
@@ -52,7 +52,7 @@ public @interface ClusterTest {
String brokerListener() default DEFAULT_BROKER_LISTENER_NAME;
SecurityProtocol controllerSecurityProtocol() default
SecurityProtocol.PLAINTEXT;
String controllerListener() default DEFAULT_CONTROLLER_LISTENER_NAME;
- MetadataVersion metadataVersion() default MetadataVersion.IBP_4_4_IV1;
+ MetadataVersion metadataVersion() default MetadataVersion.IBP_4_4_IV2;
ClusterConfigProperty[] serverProperties() default {};
// users can add tags that they want to display in test
String[] tags() default {};
diff --git
a/test-common/test-common-runtime/src/main/java/org/apache/kafka/common/test/MockController.java
b/test-common/test-common-runtime/src/main/java/org/apache/kafka/common/test/MockController.java
index cb883971b1b..9a25615d25e 100644
---
a/test-common/test-common-runtime/src/main/java/org/apache/kafka/common/test/MockController.java
+++
b/test-common/test-common-runtime/src/main/java/org/apache/kafka/common/test/MockController.java
@@ -234,6 +234,14 @@ public class MockController implements Controller {
throw new UnsupportedOperationException();
}
+ @Override
+ public CompletableFuture<Void> unregisterController(
+ ControllerRequestContext context,
+ int controllerId
+ ) {
+ throw new UnsupportedOperationException();
+ }
+
static class MockTopic {
private final String name;
private final Uuid id;
diff --git a/tools/src/main/java/org/apache/kafka/tools/ClusterTool.java
b/tools/src/main/java/org/apache/kafka/tools/ClusterTool.java
index 4342bb07bf3..0de7b10a407 100644
--- a/tools/src/main/java/org/apache/kafka/tools/ClusterTool.java
+++ b/tools/src/main/java/org/apache/kafka/tools/ClusterTool.java
@@ -79,11 +79,13 @@ public class ClusterTool {
.help("Get information about the ID of a cluster.");
Subparser unregisterParser = subparsers.addParser("unregister")
.help("Unregister a broker.");
+ Subparser unregisterControllerParser =
subparsers.addParser("unregister-controller")
+ .help("Unregister a controller.");
Subparser listEndpoints = subparsers.addParser("list-endpoints")
.help("List endpoints");
Subparser apiVersionsParser = subparsers.addParser("api-versions")
.help("Get information about the api versions of the brokers
or controllers.");
- for (Subparser subpparser : List.of(clusterIdParser, unregisterParser,
listEndpoints, apiVersionsParser)) {
+ for (Subparser subpparser : List.of(clusterIdParser, unregisterParser,
unregisterControllerParser, listEndpoints, apiVersionsParser)) {
MutuallyExclusiveGroup connectionOptions =
subpparser.addMutuallyExclusiveGroup().required(true);
connectionOptions.addArgument("--bootstrap-server", "-b")
.action(store())
@@ -104,6 +106,11 @@ public class ClusterTool {
.action(store())
.required(true)
.help("The ID of the broker to unregister.");
+ unregisterControllerParser.addArgument("--id", "-i")
+ .type(Integer.class)
+ .action(store())
+ .required(true)
+ .help("The ID of the controller to unregister.");
listEndpoints.addArgument("--include-fenced-brokers")
.action(storeTrue())
.help("Whether to include fenced brokers when listing broker
endpoints");
@@ -138,6 +145,12 @@ public class ClusterTool {
}
break;
}
+ case "unregister-controller": {
+ try (Admin adminClient = Admin.create(properties)) {
+ unregisterControllerCommand(System.out, adminClient,
namespace.getInt("id"));
+ }
+ break;
+ }
case "list-endpoints": {
try (Admin adminClient = Admin.create(properties)) {
boolean includeFencedBrokers =
Optional.of(namespace.getBoolean("include_fenced_brokers")).orElse(false);
@@ -183,6 +196,20 @@ public class ClusterTool {
}
}
+ static void unregisterControllerCommand(PrintStream stream, Admin
adminClient, int id) throws Exception {
+ try {
+ adminClient.unregisterController(id).all().get();
+ stream.println("Controller " + id + " is no longer registered.");
+ } catch (ExecutionException ee) {
+ Throwable cause = ee.getCause();
+ if (cause instanceof UnsupportedVersionException) {
+ stream.println("The target cluster does not support the
controller unregistration API.");
+ } else {
+ throw ee;
+ }
+ }
+ }
+
static void listEndpoints(PrintStream stream, Admin adminClient, boolean
listControllerEndpoints, boolean includeFencedBrokers) throws Exception {
try {
DescribeClusterOptions option = new
DescribeClusterOptions().includeFencedBrokers(includeFencedBrokers);
diff --git
a/tools/src/main/java/org/apache/kafka/tools/MetadataQuorumCommand.java
b/tools/src/main/java/org/apache/kafka/tools/MetadataQuorumCommand.java
index 6c747250bc6..484a46ce0c2 100644
--- a/tools/src/main/java/org/apache/kafka/tools/MetadataQuorumCommand.java
+++ b/tools/src/main/java/org/apache/kafka/tools/MetadataQuorumCommand.java
@@ -22,6 +22,8 @@ import org.apache.kafka.clients.admin.RaftVoterEndpoint;
import org.apache.kafka.common.Endpoint;
import org.apache.kafka.common.KafkaException;
import org.apache.kafka.common.Uuid;
+import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.errors.VoterNotFoundException;
import org.apache.kafka.common.network.ListenerName;
import org.apache.kafka.common.security.auth.SecurityProtocol;
import org.apache.kafka.common.utils.Utils;
@@ -154,7 +156,8 @@ public class MetadataQuorumCommand {
case "remove-controller" -> handleRemoveController(admin,
namespace.getInt("controller_id"),
namespace.getString("controller_directory_id"),
- namespace.getBoolean("dry_run"));
+ namespace.getBoolean("dry_run"),
+ namespace.getBoolean("unregister"));
default -> throw new IllegalStateException(format("Unknown
command: %s", command));
}
} finally {
@@ -476,13 +479,19 @@ public class MetadataQuorumCommand {
.addArgument("--dry-run")
.help("True if we should print what would be done, but not do it.")
.action(Arguments.storeTrue());
+
+ removeControllerParser
+ .addArgument("--unregister")
+ .help("If set, also unregister the controller after successfully
removing it from the voter set.")
+ .action(Arguments.storeTrue());
}
static void handleRemoveController(
Admin admin,
int controllerId,
String controllerDirectoryIdString,
- boolean dryRun
+ boolean dryRun,
+ boolean unregister
) throws TerseException, ExecutionException, InterruptedException {
if (controllerId < 0) {
throw new TerseException("Invalid negative --controller-id: " +
controllerId);
@@ -494,12 +503,55 @@ public class MetadataQuorumCommand {
throw new TerseException("Failed to parse
--controller-directory-id: " + e.getMessage());
}
if (!dryRun) {
- admin.removeRaftVoter(controllerId, directoryId).
- all().get();
+ removeRaftVoter(admin, controllerId, directoryId, unregister);
}
- System.out.printf("%s KRaft controller %d with directory id %s%n",
+ System.out.printf("%sKRaft controller %d with directory id %s%n",
dryRun ? "DRY RUN of removing " : "Removed ",
controllerId,
directoryId);
+ if (unregister) {
+ if (!dryRun) {
+ unregisterController(admin, controllerId);
+ }
+ System.out.printf("%sKRaft controller %d%n",
+ dryRun ? "DRY RUN of unregistering " : "Unregistered ",
+ controllerId);
+ }
+ }
+
+ private static void removeRaftVoter(
+ Admin admin,
+ int controllerId,
+ Uuid directoryId,
+ boolean unregister
+ ) throws TerseException, ExecutionException, InterruptedException {
+ try {
+ admin.removeRaftVoter(controllerId, directoryId).all().get();
+ } catch (ExecutionException e) {
+ Throwable cause = e.getCause();
+ if (unregister && (cause instanceof UnsupportedVersionException ||
+ cause instanceof VoterNotFoundException)) {
+ throw new TerseException("Failed to remove KRaft voter " +
controllerId
+ + ": " + cause.getMessage()
+ + ". To unregister the controller from the cluster, run "
+ + "`kafka-cluster.sh unregister-controller --id "
+ + controllerId + "`.");
+ }
+ throw e;
+ }
+ }
+
+ private static void unregisterController(Admin admin, int controllerId)
+ throws TerseException, InterruptedException {
+ try {
+ admin.unregisterController(controllerId).all().get();
+ } catch (ExecutionException e) {
+ Throwable cause = e.getCause();
+ throw new TerseException("Failed to unregister controller " +
controllerId
+ + ": " + (cause != null ? cause.getMessage() : e.getMessage())
+ + ". To unregister the controller from the cluster, run "
+ + "`kafka-cluster.sh unregister-controller --id "
+ + controllerId + "`.");
+ }
}
}
diff --git a/tools/src/test/java/org/apache/kafka/tools/ClusterToolTest.java
b/tools/src/test/java/org/apache/kafka/tools/ClusterToolTest.java
index bbd5f2d1c23..7bac04a04cc 100644
--- a/tools/src/test/java/org/apache/kafka/tools/ClusterToolTest.java
+++ b/tools/src/test/java/org/apache/kafka/tools/ClusterToolTest.java
@@ -302,6 +302,26 @@ public class ClusterToolTest {
assertEquals("Broker 0 is no longer registered.\n", stream.toString());
}
+ @Test
+ public void testUnregisterController() throws Exception {
+ Admin adminClient = new MockAdminClient.Builder().numBrokers(3).
+ usingRaftController(true).
+ build();
+ ByteArrayOutputStream stream = new ByteArrayOutputStream();
+ ClusterTool.unregisterControllerCommand(new PrintStream(stream),
adminClient, 0);
+ assertEquals("Controller 0 is no longer registered.\n",
stream.toString());
+ }
+
+ @Test
+ public void testLegacyModeClusterCannotUnregisterController() throws
Exception {
+ Admin adminClient = new MockAdminClient.Builder().numBrokers(3).
+ usingRaftController(false).
+ build();
+ ByteArrayOutputStream stream = new ByteArrayOutputStream();
+ ClusterTool.unregisterControllerCommand(new PrintStream(stream),
adminClient, 0);
+ assertEquals("The target cluster does not support the controller
unregistration API.\n", stream.toString());
+ }
+
@Test
public void testLegacyModeClusterCannotUnregisterBroker() throws Exception
{
Admin adminClient = new MockAdminClient.Builder().numBrokers(3).
diff --git
a/tools/src/test/java/org/apache/kafka/tools/MetadataQuorumCommandUnitTest.java
b/tools/src/test/java/org/apache/kafka/tools/MetadataQuorumCommandUnitTest.java
index b3131ec3b59..32d2e213306 100644
---
a/tools/src/test/java/org/apache/kafka/tools/MetadataQuorumCommandUnitTest.java
+++
b/tools/src/test/java/org/apache/kafka/tools/MetadataQuorumCommandUnitTest.java
@@ -16,8 +16,15 @@
*/
package org.apache.kafka.tools;
+import org.apache.kafka.clients.admin.Admin;
import org.apache.kafka.clients.admin.RaftVoterEndpoint;
+import org.apache.kafka.clients.admin.RemoveRaftVoterResult;
+import org.apache.kafka.clients.admin.UnregisterControllerResult;
import org.apache.kafka.common.Uuid;
+import org.apache.kafka.common.errors.UnknownServerException;
+import org.apache.kafka.common.errors.UnsupportedVersionException;
+import org.apache.kafka.common.errors.VoterNotFoundException;
+import org.apache.kafka.common.internals.KafkaFutureImpl;
import org.apache.kafka.common.utils.Utils;
import org.apache.kafka.metadata.properties.MetaProperties;
import org.apache.kafka.metadata.properties.MetaPropertiesEnsemble;
@@ -33,9 +40,12 @@ import java.util.Optional;
import java.util.Properties;
import java.util.Set;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
public class MetadataQuorumCommandUnitTest {
@Test
@@ -47,10 +57,26 @@ public class MetadataQuorumCommandUnitTest {
"--controller-id", "2",
"--controller-directory-id", "_KWDkTahTVaiVVVTaugNew",
"--dry-run"))).split("\n"));
- assertTrue(outputs.contains("DRY RUN of removing KRaft controller 2
with directory id _KWDkTahTVaiVVVTaugNew"),
+ assertTrue(outputs.contains("DRY RUN of removing KRaft controller 2
with directory id _KWDkTahTVaiVVVTaugNew"),
"Failed to find expected output in stdout: " + outputs);
}
+ @Test
+ public void testRemoveControllerDryRunWithUnregister() {
+ List<String> outputs = List.of(
+ ToolsTestUtils.captureStandardOut(() ->
+ assertEquals(0,
MetadataQuorumCommand.mainNoExit("--bootstrap-server", "localhost:9092",
+ "remove-controller",
+ "--controller-id", "2",
+ "--controller-directory-id", "_KWDkTahTVaiVVVTaugNew",
+ "--dry-run",
+ "--unregister"))).split("\n"));
+ assertTrue(outputs.contains("DRY RUN of removing KRaft controller 2
with directory id _KWDkTahTVaiVVVTaugNew"),
+ "Failed to find expected output in stdout: " + outputs);
+ assertTrue(outputs.contains("DRY RUN of unregistering KRaft controller
2"),
+ "Failed to find expected unregister output in stdout: " + outputs);
+ }
+
@Test
public void testGetControllerIdWithoutId() {
Properties props = new Properties();
@@ -260,4 +286,94 @@ public class MetadataQuorumCommandUnitTest {
"Failed to find expected output in stdout: " + outputs);
}
}
+
+ private static final int REMOVE_CONTROLLER_ID = 2;
+ private static final String REMOVE_DIRECTORY_ID_STRING =
"_KWDkTahTVaiVVVTaugNew";
+ private static final Uuid REMOVE_DIRECTORY_ID =
Uuid.fromString(REMOVE_DIRECTORY_ID_STRING);
+
+ private static <T> KafkaFutureImpl<T> completedFuture(T value) {
+ KafkaFutureImpl<T> future = new KafkaFutureImpl<>();
+ future.complete(value);
+ return future;
+ }
+
+ private static <T> KafkaFutureImpl<T> failedFuture(Throwable cause) {
+ KafkaFutureImpl<T> future = new KafkaFutureImpl<>();
+ future.completeExceptionally(cause);
+ return future;
+ }
+
+ private static Admin mockAdminForRemoveController(
+ KafkaFutureImpl<Void> removeRaftVoterFuture,
+ KafkaFutureImpl<Void> unregisterControllerFuture
+ ) {
+ Admin admin = mock(Admin.class);
+ RemoveRaftVoterResult removeResult = mock(RemoveRaftVoterResult.class);
+ when(removeResult.all()).thenReturn(removeRaftVoterFuture);
+ when(admin.removeRaftVoter(REMOVE_CONTROLLER_ID,
REMOVE_DIRECTORY_ID)).thenReturn(removeResult);
+ if (unregisterControllerFuture != null) {
+ UnregisterControllerResult unregisterResult =
mock(UnregisterControllerResult.class);
+
when(unregisterResult.all()).thenReturn(unregisterControllerFuture);
+
when(admin.unregisterController(REMOVE_CONTROLLER_ID)).thenReturn(unregisterResult);
+ }
+ return admin;
+ }
+
+ @Test
+ public void testRemoveControllerWithUnregisterSuccess() {
+ Admin admin = mockAdminForRemoveController(completedFuture(null),
completedFuture(null));
+ String stdout = ToolsTestUtils.captureStandardOut(() ->
+ assertDoesNotThrow(() ->
+ MetadataQuorumCommand.handleRemoveController(admin,
REMOVE_CONTROLLER_ID,
+ REMOVE_DIRECTORY_ID_STRING, false, true)));
+ assertTrue(stdout.contains("Removed KRaft controller " +
REMOVE_CONTROLLER_ID
+ + " with directory id " + REMOVE_DIRECTORY_ID_STRING),
+ "Expected 'Removed' line in stdout: " + stdout);
+ assertTrue(stdout.contains("Unregistered KRaft controller " +
REMOVE_CONTROLLER_ID),
+ "Expected 'Unregistered' line in stdout: " + stdout);
+ }
+
+ @Test
+ public void testRemoveControllerFailsWithUnsupportedVersion() {
+ Admin admin = mockAdminForRemoveController(
+ failedFuture(new UnsupportedVersionException("kraft.version=0
doesn't support voter removal")),
+ null);
+ TerseException e = assertThrows(TerseException.class,
+ () -> MetadataQuorumCommand.handleRemoveController(admin,
REMOVE_CONTROLLER_ID,
+ REMOVE_DIRECTORY_ID_STRING, false, true));
+ assertTrue(e.getMessage().contains("Failed to remove KRaft voter " +
REMOVE_CONTROLLER_ID),
+ "Expected removal-failure prefix; got: " + e.getMessage());
+ assertTrue(e.getMessage().contains("kafka-cluster.sh
unregister-controller --id "
+ + REMOVE_CONTROLLER_ID),
+ "Expected kafka-cluster.sh hint; got: " + e.getMessage());
+ }
+
+ @Test
+ public void testRemoveControllerFailsWithVoterNotFound() {
+ Admin admin = mockAdminForRemoveController(
+ failedFuture(new VoterNotFoundException("voter " +
REMOVE_CONTROLLER_ID + " is not in the voter set")),
+ null);
+ TerseException e = assertThrows(TerseException.class,
+ () -> MetadataQuorumCommand.handleRemoveController(admin,
REMOVE_CONTROLLER_ID,
+ REMOVE_DIRECTORY_ID_STRING, false, true));
+ assertTrue(e.getMessage().contains("Failed to remove KRaft voter " +
REMOVE_CONTROLLER_ID),
+ "Expected removal-failure prefix; got: " + e.getMessage());
+ assertTrue(e.getMessage().contains("kafka-cluster.sh
unregister-controller --id "
+ + REMOVE_CONTROLLER_ID),
+ "Expected kafka-cluster.sh hint; got: " + e.getMessage());
+ }
+
+ @Test
+ public void testRemoveControllerSucceedsButUnregisterFails() {
+ Admin admin = mockAdminForRemoveController(completedFuture(null),
+ failedFuture(new UnknownServerException("server failed to apply
unregister")));
+ TerseException e = assertThrows(TerseException.class,
+ () -> MetadataQuorumCommand.handleRemoveController(admin,
REMOVE_CONTROLLER_ID,
+ REMOVE_DIRECTORY_ID_STRING, false, true));
+ assertTrue(e.getMessage().contains("Failed to unregister controller "
+ REMOVE_CONTROLLER_ID),
+ "Expected post-removal failure prefix; got: " + e.getMessage());
+ assertTrue(e.getMessage().contains("kafka-cluster.sh
unregister-controller --id "
+ + REMOVE_CONTROLLER_ID),
+ "Expected kafka-cluster.sh hint; got: " + e.getMessage());
+ }
}