From 27c0c40111fa1a1a751ae242211b3cb98aa118f9 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 2 Jun 2022 08:58:36 +0200 Subject: [PATCH 01/18] add new metadata version --- .../org/apache/kafka/server/common/MetadataVersion.java | 9 ++++++++- .../apache/kafka/server/common/MetadataVersionTest.java | 6 +++++- 2 files changed, 13 insertions(+), 2 deletions(-) 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 ee19faf88b54b..cdf69f993cbd5 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 @@ -159,7 +159,10 @@ public enum MetadataVersion { IBP_3_3_IV1(6, "3.3", "IV1", true), // In KRaft mode, use BrokerRegistrationChangeRecord instead of UnfenceBrokerRecord and FenceBrokerRecord. - IBP_3_3_IV2(7, "3.3", "IV2", true); + IBP_3_3_IV2(7, "3.3", "IV2", true), + + // Adds InControlledShutdown state to RegisterBrokerRecord and BrokerRegistrationChangeRecord (KIP-841). + IBP_3_3_IV3(8, "3.3", "IV3", true); public static final String FEATURE_NAME = "metadata.version"; @@ -243,6 +246,10 @@ public boolean isBrokerRegistrationChangeRecordSupported() { return this.isAtLeast(IBP_3_3_IV2); } + public boolean isInControlledShutdownStateSupported() { + return this.isAtLeast(IBP_3_3_IV3); + } + private static final Map IBP_VERSIONS; static { { diff --git a/server-common/src/test/java/org/apache/kafka/server/common/MetadataVersionTest.java b/server-common/src/test/java/org/apache/kafka/server/common/MetadataVersionTest.java index 07a644e729396..ec038383caad4 100644 --- a/server-common/src/test/java/org/apache/kafka/server/common/MetadataVersionTest.java +++ b/server-common/src/test/java/org/apache/kafka/server/common/MetadataVersionTest.java @@ -62,6 +62,7 @@ import static org.apache.kafka.server.common.MetadataVersion.IBP_3_3_IV0; import static org.apache.kafka.server.common.MetadataVersion.IBP_3_3_IV1; import static org.apache.kafka.server.common.MetadataVersion.IBP_3_3_IV2; +import static org.apache.kafka.server.common.MetadataVersion.IBP_3_3_IV3; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; @@ -190,10 +191,11 @@ public void testFromVersionString() { assertEquals(IBP_3_2_IV0, MetadataVersion.fromVersionString("3.2")); assertEquals(IBP_3_2_IV0, MetadataVersion.fromVersionString("3.2-IV0")); - assertEquals(IBP_3_3_IV2, MetadataVersion.fromVersionString("3.3")); + assertEquals(IBP_3_3_IV3, MetadataVersion.fromVersionString("3.3")); assertEquals(IBP_3_3_IV0, MetadataVersion.fromVersionString("3.3-IV0")); assertEquals(IBP_3_3_IV1, MetadataVersion.fromVersionString("3.3-IV1")); assertEquals(IBP_3_3_IV2, MetadataVersion.fromVersionString("3.3-IV2")); + assertEquals(IBP_3_3_IV3, MetadataVersion.fromVersionString("3.3-IV3")); } @Test @@ -240,6 +242,7 @@ public void testShortVersion() { assertEquals("3.3", IBP_3_3_IV0.shortVersion()); assertEquals("3.3", IBP_3_3_IV1.shortVersion()); assertEquals("3.3", IBP_3_3_IV2.shortVersion()); + assertEquals("3.3", IBP_3_3_IV3.shortVersion()); } @Test @@ -275,6 +278,7 @@ public void testVersion() { assertEquals("3.3-IV0", IBP_3_3_IV0.version()); assertEquals("3.3-IV1", IBP_3_3_IV1.version()); assertEquals("3.3-IV2", IBP_3_3_IV2.version()); + assertEquals("3.3-IV3", IBP_3_3_IV3.version()); } @Test From 4a1052cfeaf74436b5c6be3188f51ba7231e6ee6 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 2 Jun 2022 09:11:12 +0200 Subject: [PATCH 02/18] update BrokerRegistrationRecord and BrokerRegistrationChangeRecord --- .../org/apache/kafka/controller/ClusterControlManager.java | 5 ++--- .../java/org/apache/kafka/metadata/BrokerRegistration.java | 6 +----- .../common/metadata/BrokerRegistrationChangeRecord.json | 6 ++++-- .../resources/common/metadata/RegisterBrokerRecord.json | 6 ++++-- 4 files changed, 11 insertions(+), 12 deletions(-) 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 5c069e12e7e3a..7247c105fc3f5 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -65,7 +65,6 @@ import java.util.stream.Collectors; import static java.util.concurrent.TimeUnit.NANOSECONDS; -import static org.apache.kafka.common.metadata.MetadataRecordType.REGISTER_BROKER_RECORD; /** @@ -339,7 +338,7 @@ public ControllerResult registerBroker( heartbeatManager.register(brokerId, record.fenced()); List records = new ArrayList<>(); - records.add(new ApiMessageAndVersion(record, REGISTER_BROKER_RECORD.highestSupportedVersion())); + records.add(new ApiMessageAndVersion(record, (short) 0)); return ControllerResult.atomicOf(records, new BrokerRegistrationReply(brokerEpoch)); } @@ -557,7 +556,7 @@ public List next() { setEndPoints(endpoints). setFeatures(features). setRack(registration.rack().orElse(null)). - setFenced(registration.fenced()), REGISTER_BROKER_RECORD.highestSupportedVersion())); + setFenced(registration.fenced()), (short) 0)); return batch; } } diff --git a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java index cc8ed9b4aaebb..8602d824aed1a 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java @@ -36,9 +36,6 @@ import java.util.Optional; import java.util.stream.Collectors; -import static org.apache.kafka.common.metadata.MetadataRecordType.REGISTER_BROKER_RECORD; - - /** * An immutable class which represents broker registrations. */ @@ -173,8 +170,7 @@ public ApiMessageAndVersion toRecord() { setMinSupportedVersion(entry.getValue().min()). setMaxSupportedVersion(entry.getValue().max())); } - return new ApiMessageAndVersion(registrationRecord, - REGISTER_BROKER_RECORD.highestSupportedVersion()); + return new ApiMessageAndVersion(registrationRecord, (short) 0); } @Override diff --git a/metadata/src/main/resources/common/metadata/BrokerRegistrationChangeRecord.json b/metadata/src/main/resources/common/metadata/BrokerRegistrationChangeRecord.json index 152508ce54f09..81bebaaff276c 100644 --- a/metadata/src/main/resources/common/metadata/BrokerRegistrationChangeRecord.json +++ b/metadata/src/main/resources/common/metadata/BrokerRegistrationChangeRecord.json @@ -17,7 +17,7 @@ "apiKey": 17, "type": "metadata", "name": "BrokerRegistrationChangeRecord", - "validVersions": "0", + "validVersions": "0-1", "flexibleVersions": "0+", "fields": [ { "name": "BrokerId", "type": "int32", "versions": "0+", "entityType": "brokerId", @@ -25,6 +25,8 @@ { "name": "BrokerEpoch", "type": "int64", "versions": "0+", "about": "The broker epoch assigned by the controller." }, { "name": "Fenced", "type": "int8", "versions": "0+", "taggedVersions": "0+", "tag": 0, - "about": "-1 if the broker has been unfenced, 0 if no change, 1 if the broker has been fenced." } + "about": "-1 if the broker has been unfenced, 0 if no change, 1 if the broker has been fenced." }, + { "name": "InControlledShutdown", "type": "int8", "versions": "1+", "taggedVersions": "1+", "tag": 1, + "about": "0 if no change, 1 if the broker is in controlled shutdown." } ] } diff --git a/metadata/src/main/resources/common/metadata/RegisterBrokerRecord.json b/metadata/src/main/resources/common/metadata/RegisterBrokerRecord.json index a0e7af2fbed8c..a32c16d8a607c 100644 --- a/metadata/src/main/resources/common/metadata/RegisterBrokerRecord.json +++ b/metadata/src/main/resources/common/metadata/RegisterBrokerRecord.json @@ -17,7 +17,7 @@ "apiKey": 0, "type": "metadata", "name": "RegisterBrokerRecord", - "validVersions": "0", + "validVersions": "0-1", "flexibleVersions": "0+", "fields": [ { "name": "BrokerId", "type": "int32", "versions": "0+", "entityType": "brokerId", @@ -49,6 +49,8 @@ { "name": "Rack", "type": "string", "versions": "0+", "nullableVersions": "0+", "about": "The broker rack." }, { "name": "Fenced", "type": "bool", "versions": "0+", "default": "true", - "about": "True if the broker is fenced." } + "about": "True if the broker is fenced." }, + { "name": "InControlledShutdown", "type": "bool", "versions": "1+", "default": "false", + "about": "True if the broker is in controlled shutdown." } ] } From 85a76a4fb343775d30bc4fcd3d85980c55b981dc Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 2 Jun 2022 11:32:51 +0200 Subject: [PATCH 03/18] wire RegisterBrokerRecord --- .../controller/ClusterControlManager.java | 58 +++++++-- .../kafka/controller/QuorumController.java | 11 +- .../kafka/metadata/BrokerRegistration.java | 27 ++-- .../controller/ClusterControlManagerTest.java | 118 ++++++++++++++++-- .../controller/QuorumControllerTest.java | 8 +- .../apache/kafka/image/ClusterImageTest.java | 9 +- .../metadata/BrokerRegistrationTest.java | 25 ++-- 7 files changed, 209 insertions(+), 47 deletions(-) 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 7247c105fc3f5..48e69a55ada5d 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -17,6 +17,7 @@ package org.apache.kafka.controller; +import org.apache.kafka.clients.ApiVersions; import org.apache.kafka.common.Endpoint; import org.apache.kafka.common.Uuid; import org.apache.kafka.common.errors.DuplicateBrokerRegistrationException; @@ -46,11 +47,13 @@ import org.apache.kafka.metadata.placement.StripedReplicaPlacer; import org.apache.kafka.metadata.placement.UsableBroker; import org.apache.kafka.server.common.ApiMessageAndVersion; +import org.apache.kafka.server.common.MetadataVersion; import org.apache.kafka.timeline.SnapshotRegistry; import org.apache.kafka.timeline.TimelineHashMap; import org.slf4j.Logger; import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; import java.util.Iterator; import java.util.List; @@ -83,6 +86,7 @@ static class Builder { private long sessionTimeoutNs = DEFAULT_SESSION_TIMEOUT_NS; private ReplicaPlacer replicaPlacer = null; private ControllerMetrics controllerMetrics = null; + private FeatureControlManager featureControl = null; Builder setLogContext(LogContext logContext) { this.logContext = logContext; @@ -119,8 +123,15 @@ Builder setControllerMetrics(ControllerMetrics controllerMetrics) { return this; } + Builder setFeatureControlManager(FeatureControlManager featureControl) { + this.featureControl = featureControl; + return this; + } + ClusterControlManager build() { - if (logContext == null) logContext = new LogContext(); + if (logContext == null) { + logContext = new LogContext(); + } if (clusterId == null) { clusterId = Uuid.randomUuid().toString(); } @@ -131,7 +142,17 @@ ClusterControlManager build() { replicaPlacer = new StripedReplicaPlacer(new Random()); } if (controllerMetrics == null) { - throw new RuntimeException("You must specify controllerMetrics"); + throw new RuntimeException("You must specify ControllerMetrics"); + } + if (featureControl == null) { + featureControl = new FeatureControlManager.Builder(). + setLogContext(logContext). + setSnapshotRegistry(snapshotRegistry). + setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), + QuorumFeatures.defaultFeatureMap(), + Collections.singletonList(0))). + setMetadataVersion(MetadataVersion.latest()). + build(); } return new ClusterControlManager(logContext, clusterId, @@ -139,7 +160,9 @@ ClusterControlManager build() { snapshotRegistry, sessionTimeoutNs, replicaPlacer, - controllerMetrics); + controllerMetrics, + featureControl + ); } } @@ -217,6 +240,11 @@ boolean check() { */ private Optional readyBrokersFuture; + /** + * The feature control manager. + */ + private final FeatureControlManager featureControl; + private ClusterControlManager( LogContext logContext, String clusterId, @@ -224,7 +252,8 @@ private ClusterControlManager( SnapshotRegistry snapshotRegistry, long sessionTimeoutNs, ReplicaPlacer replicaPlacer, - ControllerMetrics metrics + ControllerMetrics metrics, + FeatureControlManager featureControl ) { this.logContext = logContext; this.clusterId = clusterId; @@ -236,6 +265,7 @@ private ClusterControlManager( this.heartbeatManager = null; this.readyBrokersFuture = Optional.empty(); this.controllerMetrics = metrics; + this.featureControl = featureControl; } ReplicaPlacer replicaPlacer() { @@ -278,6 +308,13 @@ Set fencedBrokerIds() { .collect(Collectors.toSet()); } + private short registerBrokerRecordVersion() { + if (featureControl.metadataVersion().isInControlledShutdownStateSupported()) + return (short) 1; + else + return (short) 0; + } + /** * Process an incoming broker registration request. */ @@ -338,7 +375,7 @@ public ControllerResult registerBroker( heartbeatManager.register(brokerId, record.fenced()); List records = new ArrayList<>(); - records.add(new ApiMessageAndVersion(record, (short) 0)); + records.add(new ApiMessageAndVersion(record, registerBrokerRecordVersion())); return ControllerResult.atomicOf(records, new BrokerRegistrationReply(brokerEpoch)); } @@ -360,7 +397,8 @@ public void replay(RegisterBrokerRecord record) { BrokerRegistration prevRegistration = brokerRegistrations.put(brokerId, new BrokerRegistration(brokerId, record.brokerEpoch(), record.incarnationId(), listeners, features, - Optional.ofNullable(record.rack()), record.fenced())); + Optional.ofNullable(record.rack()), record.fenced(), + record.inControlledShutdown())); updateMetrics(prevRegistration, brokerRegistrations.get(brokerId)); if (heartbeatManager != null) { if (prevRegistration != null) heartbeatManager.remove(brokerId); @@ -490,6 +528,12 @@ public boolean unfenced(int brokerId) { return !registration.fenced(); } + public boolean inControlledShutdown(int brokerId) { + BrokerRegistration registration = brokerRegistrations.get(brokerId); + if (registration == null) return false; + return registration.inControlledShutdown(); + } + BrokerHeartbeatManager heartbeatManager() { if (heartbeatManager == null) { throw new RuntimeException("ClusterControlManager is not active."); @@ -556,7 +600,7 @@ public List next() { setEndPoints(endpoints). setFeatures(features). setRack(registration.rack().orElse(null)). - setFenced(registration.fenced()), (short) 0)); + setFenced(registration.fenced()), registerBrokerRecordVersion())); return batch; } } 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 e921ac4067ada..ea9b7205e5050 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -1557,6 +1557,11 @@ private QuorumController(LogContext logContext, setNodeId(nodeId). build(); this.clientQuotaControlManager = new ClientQuotaControlManager(snapshotRegistry); + this.featureControl = new FeatureControlManager.Builder(). + setLogContext(logContext). + setQuorumFeatures(quorumFeatures). + setSnapshotRegistry(snapshotRegistry). + build(); this.clusterControl = new ClusterControlManager.Builder(). setLogContext(logContext). setClusterId(clusterId). @@ -1565,12 +1570,8 @@ private QuorumController(LogContext logContext, setSessionTimeoutNs(sessionTimeoutNs). setReplicaPlacer(replicaPlacer). setControllerMetrics(controllerMetrics). + setFeatureControlManager(featureControl). build(); - this.featureControl = new FeatureControlManager.Builder(). - setLogContext(logContext). - setQuorumFeatures(quorumFeatures). - setSnapshotRegistry(snapshotRegistry). - build(); this.producerIdControlManager = new ProducerIdControlManager(clusterControl, snapshotRegistry); this.snapshotMaxNewRecordBytes = snapshotMaxNewRecordBytes; this.leaderImbalanceCheckIntervalNs = leaderImbalanceCheckIntervalNs; diff --git a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java index 8602d824aed1a..155eec557b6f9 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java @@ -55,6 +55,7 @@ private static Map listenersToMap(Collection listene private final Map supportedFeatures; private final Optional rack; private final boolean fenced; + private final boolean inControlledShutdown; public BrokerRegistration(int id, long epoch, @@ -62,8 +63,10 @@ public BrokerRegistration(int id, List listeners, Map supportedFeatures, Optional rack, - boolean fenced) { - this(id, epoch, incarnationId, listenersToMap(listeners), supportedFeatures, rack, fenced); + boolean fenced, + boolean inControlledShutdown) { + this(id, epoch, incarnationId, listenersToMap(listeners), supportedFeatures, rack, + fenced, inControlledShutdown); } public BrokerRegistration(int id, @@ -72,7 +75,8 @@ public BrokerRegistration(int id, Map listeners, Map supportedFeatures, Optional rack, - boolean fenced) { + boolean fenced, + boolean inControlledShutdown) { this.id = id; this.epoch = epoch; this.incarnationId = incarnationId; @@ -89,6 +93,7 @@ public BrokerRegistration(int id, Objects.requireNonNull(rack); this.rack = rack; this.fenced = fenced; + this.inControlledShutdown = inControlledShutdown; } public static BrokerRegistration fromRecord(RegisterBrokerRecord record) { @@ -110,7 +115,8 @@ public static BrokerRegistration fromRecord(RegisterBrokerRecord record) { listeners, supportedFeatures, Optional.ofNullable(record.rack()), - record.fenced()); + record.fenced(), + record.inControlledShutdown()); } public int id() { @@ -149,13 +155,18 @@ public boolean fenced() { return fenced; } + public boolean inControlledShutdown() { + return inControlledShutdown; + } + public ApiMessageAndVersion toRecord() { RegisterBrokerRecord registrationRecord = new RegisterBrokerRecord(). setBrokerId(id). setRack(rack.orElse(null)). setBrokerEpoch(epoch). setIncarnationId(incarnationId). - setFenced(fenced); + setFenced(fenced). + setInControlledShutdown(inControlledShutdown); for (Entry entry : listeners.entrySet()) { Endpoint endpoint = entry.getValue(); registrationRecord.endPoints().add(new BrokerEndpoint(). @@ -189,7 +200,8 @@ public boolean equals(Object o) { other.listeners.equals(listeners) && other.supportedFeatures.equals(supportedFeatures) && other.rack.equals(rack) && - other.fenced == fenced; + other.fenced == fenced && + other.inControlledShutdown == inControlledShutdown; } @Override @@ -209,12 +221,13 @@ public String toString() { bld.append("}"); bld.append(", rack=").append(rack); bld.append(", fenced=").append(fenced); + bld.append(", inControlledShutdown=").append(inControlledShutdown); bld.append(")"); return bld.toString(); } public BrokerRegistration cloneWithFencing(boolean fencing) { return new BrokerRegistration(id, epoch, incarnationId, listeners, - supportedFeatures, rack, fencing); + supportedFeatures, rack, fencing, inControlledShutdown); } } 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 590578f6c63bb..b5b5d2f4b2193 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ClusterControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ClusterControlManagerTest.java @@ -24,6 +24,7 @@ import java.util.List; import java.util.Optional; +import org.apache.kafka.clients.ApiVersions; import org.apache.kafka.common.Endpoint; import org.apache.kafka.common.Uuid; import org.apache.kafka.common.errors.InconsistentClusterIdException; @@ -40,6 +41,7 @@ import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.common.utils.MockTime; import org.apache.kafka.metadata.BrokerRegistration; +import org.apache.kafka.metadata.BrokerRegistrationReply; import org.apache.kafka.metadata.FinalizedControllerFeatures; import org.apache.kafka.metadata.RecordTestUtils; import org.apache.kafka.metadata.placement.ClusterDescriber; @@ -55,6 +57,7 @@ import org.junit.jupiter.params.provider.ValueSource; import static org.apache.kafka.server.common.MetadataVersion.IBP_3_3_IV2; +import static org.apache.kafka.server.common.MetadataVersion.IBP_3_3_IV3; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; @@ -118,6 +121,50 @@ public void testReplay(MetadataVersion metadataVersion) { assertFalse(clusterControl.unfenced(1)); } + public void testReplayRegisterBrokerRecord() { + MockTime time = new MockTime(0, 0, 0); + + SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); + ClusterControlManager clusterControl = new ClusterControlManager.Builder(). + setClusterId("fPZv1VBsRFmnlRvmGcOW9w"). + setTime(time). + setSnapshotRegistry(snapshotRegistry). + setSessionTimeoutNs(1000). + setControllerMetrics(new MockControllerMetrics()). + build(); + + assertFalse(clusterControl.unfenced(0)); + assertFalse(clusterControl.inControlledShutdown(0)); + + RegisterBrokerRecord brokerRecord = new RegisterBrokerRecord(). + setBrokerEpoch(100). + setBrokerId(0). + setRack(null). + setFenced(true). + setInControlledShutdown(true); + brokerRecord.endPoints().add(new BrokerEndpoint(). + setSecurityProtocol(SecurityProtocol.PLAINTEXT.id). + setPort((short) 9092). + setName("PLAINTEXT"). + setHost("example.com")); + clusterControl.replay(brokerRecord); + + assertFalse(clusterControl.unfenced(0)); + assertTrue(clusterControl.inControlledShutdown(0)); + + brokerRecord.setInControlledShutdown(false); + clusterControl.replay(brokerRecord); + + assertFalse(clusterControl.unfenced(0)); + assertFalse(clusterControl.inControlledShutdown(0)); + + brokerRecord.setFenced(false); + clusterControl.replay(brokerRecord); + + assertTrue(clusterControl.unfenced(0)); + assertFalse(clusterControl.inControlledShutdown(0)); + } + @Test public void testRegistrationWithIncorrectClusterId() throws Exception { SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); @@ -139,6 +186,49 @@ public void testRegistrationWithIncorrectClusterId() throws Exception { new FinalizedControllerFeatures(Collections.emptyMap(), 456L))); } + @ParameterizedTest + @EnumSource(value = MetadataVersion.class, names = {"IBP_3_3_IV2", "IBP_3_3_IV3"}) + public void testRegisterBrokerRecordVersion(MetadataVersion metadataVersion) { + SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); + FeatureControlManager featureControl = new FeatureControlManager.Builder(). + setSnapshotRegistry(snapshotRegistry). + setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), + QuorumFeatures.defaultFeatureMap(), + Collections.singletonList(0))). + setMetadataVersion(metadataVersion). + build(); + ClusterControlManager clusterControl = new ClusterControlManager.Builder(). + setClusterId("fPZv1VBsRFmnlRvmGcOW9w"). + setTime(new MockTime(0, 0, 0)). + setSnapshotRegistry(snapshotRegistry). + setSessionTimeoutNs(1000). + setControllerMetrics(new MockControllerMetrics()). + setFeatureControlManager(featureControl). + build(); + clusterControl.activate(); + + ControllerResult result = clusterControl.registerBroker( + new BrokerRegistrationRequestData(). + setClusterId("fPZv1VBsRFmnlRvmGcOW9w"). + setBrokerId(0). + setRack(null). + setIncarnationId(Uuid.fromString("0H4fUu1xQEKXFYwB1aBjhg")), + 123L, + new FinalizedControllerFeatures(Collections.emptyMap(), 456L)); + + short expectedVersion = metadataVersion.isAtLeast(IBP_3_3_IV3) ? (short) 1 : (short) 0; + + assertEquals( + Arrays.asList(new ApiMessageAndVersion(new RegisterBrokerRecord(). + setBrokerEpoch(123L). + setBrokerId(0). + setRack(null). + setIncarnationId(Uuid.fromString("0H4fUu1xQEKXFYwB1aBjhg")). + setFenced(true). + setInControlledShutdown(false), expectedVersion)), + result.records()); + } + @Test public void testUnregister() throws Exception { RegisterBrokerRecord brokerRecord = new RegisterBrokerRecord(). @@ -161,10 +251,10 @@ public void testUnregister() throws Exception { clusterControl.activate(); clusterControl.replay(brokerRecord); assertEquals(new BrokerRegistration(1, 100, - Uuid.fromString("fPZv1VBsRFmnlRvmGcOW9w"), Collections.singletonMap("PLAINTEXT", - new Endpoint("PLAINTEXT", SecurityProtocol.PLAINTEXT, "example.com", 9092)), - Collections.emptyMap(), Optional.of("arack"), true), - clusterControl.brokerRegistrations().get(1)); + Uuid.fromString("fPZv1VBsRFmnlRvmGcOW9w"), Collections.singletonMap("PLAINTEXT", + new Endpoint("PLAINTEXT", SecurityProtocol.PLAINTEXT, "example.com", 9092)), + Collections.emptyMap(), Optional.of("arack"), true, false), + clusterControl.brokerRegistrations().get(1)); UnregisterBrokerRecord unregisterRecord = new UnregisterBrokerRecord(). setBrokerId(1). setBrokerEpoch(100); @@ -223,15 +313,24 @@ public Iterator usableBrokers() { } } - @Test - public void testIterator() throws Exception { + @ParameterizedTest + @EnumSource(value = MetadataVersion.class, names = {"IBP_3_3_IV2", "IBP_3_3_IV3"}) + public void testIterator(MetadataVersion metadataVersion) throws Exception { MockTime time = new MockTime(0, 0, 0); SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); + FeatureControlManager featureControl = new FeatureControlManager.Builder(). + setSnapshotRegistry(snapshotRegistry). + setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), + QuorumFeatures.defaultFeatureMap(), + Collections.singletonList(0))). + setMetadataVersion(metadataVersion). + build(); ClusterControlManager clusterControl = new ClusterControlManager.Builder(). setTime(time). setSnapshotRegistry(snapshotRegistry). setSessionTimeoutNs(1000). setControllerMetrics(new MockControllerMetrics()). + setFeatureControlManager(featureControl). build(); clusterControl.activate(); assertFalse(clusterControl.unfenced(0)); @@ -250,6 +349,7 @@ public void testIterator() throws Exception { new UnfenceBrokerRecord().setId(i).setEpoch(100); clusterControl.replay(unfenceBrokerRecord); } + short expectedVersion = metadataVersion.isAtLeast(IBP_3_3_IV3) ? (short) 1 : (short) 0; RecordTestUtils.assertBatchIteratorContains(Arrays.asList( Arrays.asList(new ApiMessageAndVersion(new RegisterBrokerRecord(). setBrokerEpoch(100).setBrokerId(0).setRack(null). @@ -258,7 +358,7 @@ public void testIterator() throws Exception { setPort((short) 9092). setName("PLAINTEXT"). setHost("example.com")).iterator())). - setFenced(false), (short) 0)), + setFenced(false), expectedVersion)), Arrays.asList(new ApiMessageAndVersion(new RegisterBrokerRecord(). setBrokerEpoch(100).setBrokerId(1).setRack(null). setEndPoints(new BrokerEndpointCollection(Collections.singleton( @@ -266,7 +366,7 @@ public void testIterator() throws Exception { setPort((short) 9093). setName("PLAINTEXT"). setHost("example.com")).iterator())). - setFenced(false), (short) 0)), + setFenced(false), expectedVersion)), Arrays.asList(new ApiMessageAndVersion(new RegisterBrokerRecord(). setBrokerEpoch(100).setBrokerId(2).setRack(null). setEndPoints(new BrokerEndpointCollection(Collections.singleton( @@ -274,7 +374,7 @@ public void testIterator() throws Exception { setPort((short) 9094). setName("PLAINTEXT"). setHost("example.com")).iterator())). - setFenced(true), (short) 0))), + setFenced(true), expectedVersion))), clusterControl.iterator(Long.MAX_VALUE)); } } diff --git a/metadata/src/test/java/org/apache/kafka/controller/QuorumControllerTest.java b/metadata/src/test/java/org/apache/kafka/controller/QuorumControllerTest.java index 4e898b924c80c..d429a13b0ff48 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/QuorumControllerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/QuorumControllerTest.java @@ -760,7 +760,7 @@ private List expectedSnapshotContent(Uuid fooId, Map expectedSnapshotContent(Uuid fooId, Map expectedSnapshotContent(Uuid fooId, Map Date: Thu, 2 Jun 2022 11:54:40 +0200 Subject: [PATCH 04/18] wire BrokerRegistrationChangeRecord --- .../controller/ClusterControlManager.java | 21 ++++++-- .../kafka/metadata/BrokerRegistration.java | 5 ++ ...egistrationInControlledShutdownChange.java | 54 +++++++++++++++++++ .../controller/ClusterControlManagerTest.java | 53 ++++++++++++++++++ 4 files changed, 129 insertions(+), 4 deletions(-) create mode 100644 metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistrationInControlledShutdownChange.java 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 48e69a55ada5d..c2e322781d10a 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -40,6 +40,7 @@ import org.apache.kafka.common.utils.Time; import org.apache.kafka.metadata.BrokerRegistration; import org.apache.kafka.metadata.BrokerRegistrationFencingChange; +import org.apache.kafka.metadata.BrokerRegistrationInControlledShutdownChange; import org.apache.kafka.metadata.BrokerRegistrationReply; import org.apache.kafka.metadata.FinalizedControllerFeatures; import org.apache.kafka.metadata.VersionRange; @@ -432,12 +433,14 @@ public void replay(UnregisterBrokerRecord record) { public void replay(FenceBrokerRecord record) { replayRegistrationChange(record, record.id(), record.epoch(), - BrokerRegistrationFencingChange.UNFENCE); + BrokerRegistrationFencingChange.UNFENCE, + BrokerRegistrationInControlledShutdownChange.NONE); } public void replay(UnfenceBrokerRecord record) { replayRegistrationChange(record, record.id(), record.epoch(), - BrokerRegistrationFencingChange.FENCE); + BrokerRegistrationFencingChange.FENCE, + BrokerRegistrationInControlledShutdownChange.NONE); } public void replay(BrokerRegistrationChangeRecord record) { @@ -447,15 +450,22 @@ public void replay(BrokerRegistrationChangeRecord record) { throw new RuntimeException(String.format("Unable to replay %s: unknown " + "value for fenced field: %d", record.toString(), record.fenced())); } + Optional inControlledShutdownChange = + BrokerRegistrationInControlledShutdownChange.fromValue(record.inControlledShutdown()); + if (!inControlledShutdownChange.isPresent()) { + throw new RuntimeException(String.format("Unable to replay %s: unknown " + + "value for inControlledShutdown field: %d", record.toString(), record.inControlledShutdown())); + } replayRegistrationChange(record, record.brokerId(), record.brokerEpoch(), - fencingChange.get()); + fencingChange.get(), inControlledShutdownChange.get()); } private void replayRegistrationChange( ApiMessage record, int brokerId, long brokerEpoch, - BrokerRegistrationFencingChange fencingChange + BrokerRegistrationFencingChange fencingChange, + BrokerRegistrationInControlledShutdownChange inControlledShutdownChange ) { BrokerRegistration curRegistration = brokerRegistrations.get(brokerId); if (curRegistration == null) { @@ -469,6 +479,9 @@ private void replayRegistrationChange( if (fencingChange != BrokerRegistrationFencingChange.NONE) { nextRegistration = nextRegistration.cloneWithFencing(fencingChange.asBoolean().get()); } + if (inControlledShutdownChange != BrokerRegistrationInControlledShutdownChange.NONE) { + nextRegistration = nextRegistration.cloneWithInControlledShutdown(inControlledShutdownChange.asBoolean().get()); + } if (!curRegistration.equals(nextRegistration)) { brokerRegistrations.put(brokerId, nextRegistration); updateMetrics(curRegistration, nextRegistration); diff --git a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java index 155eec557b6f9..64d967c68b378 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java @@ -230,4 +230,9 @@ public BrokerRegistration cloneWithFencing(boolean fencing) { return new BrokerRegistration(id, epoch, incarnationId, listeners, supportedFeatures, rack, fencing, inControlledShutdown); } + + public BrokerRegistration cloneWithInControlledShutdown(boolean inControlledShutdown) { + return new BrokerRegistration(id, epoch, incarnationId, listeners, + supportedFeatures, rack, fenced, inControlledShutdown); + } } diff --git a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistrationInControlledShutdownChange.java b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistrationInControlledShutdownChange.java new file mode 100644 index 0000000000000..0cd4c3959a249 --- /dev/null +++ b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistrationInControlledShutdownChange.java @@ -0,0 +1,54 @@ +/* + * 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.metadata; + +import java.util.Arrays; +import java.util.Map; +import java.util.Optional; +import java.util.function.Function; +import java.util.stream.Collectors; + +public enum BrokerRegistrationInControlledShutdownChange { + NONE(0, Optional.empty()), + IN_CONTROLLED_SHUTDOWN(1, Optional.of(true)); + + private final byte value; + + private final Optional asBoolean; + + private final static Map VALUE_TO_ENUM = + Arrays.stream(BrokerRegistrationInControlledShutdownChange.values()). + collect(Collectors.toMap(v -> Byte.valueOf(v.value()), Function.identity())); + + public static Optional fromValue(byte value) { + return Optional.ofNullable(VALUE_TO_ENUM.get(value)); + } + + BrokerRegistrationInControlledShutdownChange(int value, Optional asBoolean) { + this.value = (byte) value; + this.asBoolean = asBoolean; + } + + public Optional asBoolean() { + return asBoolean; + } + + public byte value() { + return value; + } +} 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 b5b5d2f4b2193..d5da9d921b5b4 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ClusterControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ClusterControlManagerTest.java @@ -41,6 +41,8 @@ import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.common.utils.MockTime; import org.apache.kafka.metadata.BrokerRegistration; +import org.apache.kafka.metadata.BrokerRegistrationFencingChange; +import org.apache.kafka.metadata.BrokerRegistrationInControlledShutdownChange; import org.apache.kafka.metadata.BrokerRegistrationReply; import org.apache.kafka.metadata.FinalizedControllerFeatures; import org.apache.kafka.metadata.RecordTestUtils; @@ -121,6 +123,7 @@ public void testReplay(MetadataVersion metadataVersion) { assertFalse(clusterControl.unfenced(1)); } + @Test public void testReplayRegisterBrokerRecord() { MockTime time = new MockTime(0, 0, 0); @@ -165,6 +168,56 @@ public void testReplayRegisterBrokerRecord() { assertFalse(clusterControl.inControlledShutdown(0)); } + @Test + public void testReplayBrokerRegistrationChangeRecord() { + MockTime time = new MockTime(0, 0, 0); + + SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); + ClusterControlManager clusterControl = new ClusterControlManager.Builder(). + setClusterId("fPZv1VBsRFmnlRvmGcOW9w"). + setTime(time). + setSnapshotRegistry(snapshotRegistry). + setSessionTimeoutNs(1000). + setControllerMetrics(new MockControllerMetrics()). + build(); + + assertFalse(clusterControl.unfenced(0)); + assertFalse(clusterControl.inControlledShutdown(0)); + + RegisterBrokerRecord brokerRecord = new RegisterBrokerRecord(). + setBrokerEpoch(100). + setBrokerId(0). + setRack(null). + setFenced(false); + brokerRecord.endPoints().add(new BrokerEndpoint(). + setSecurityProtocol(SecurityProtocol.PLAINTEXT.id). + setPort((short) 9092). + setName("PLAINTEXT"). + setHost("example.com")); + clusterControl.replay(brokerRecord); + + assertTrue(clusterControl.unfenced(0)); + assertFalse(clusterControl.inControlledShutdown(0)); + + BrokerRegistrationChangeRecord registrationChangeRecord = new BrokerRegistrationChangeRecord() + .setBrokerId(0) + .setBrokerEpoch(100) + .setInControlledShutdown(BrokerRegistrationInControlledShutdownChange.IN_CONTROLLED_SHUTDOWN.value()); + clusterControl.replay(registrationChangeRecord); + + assertTrue(clusterControl.unfenced(0)); + assertTrue(clusterControl.inControlledShutdown(0)); + + registrationChangeRecord = new BrokerRegistrationChangeRecord() + .setBrokerId(0) + .setBrokerEpoch(100) + .setFenced(BrokerRegistrationFencingChange.FENCE.value()); + clusterControl.replay(registrationChangeRecord); + + assertTrue(clusterControl.unfenced(0)); + assertTrue(clusterControl.inControlledShutdown(0)); + } + @Test public void testRegistrationWithIncorrectClusterId() throws Exception { SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); From 6da0fed4e3cb4c668ccb993f103a0ed0b17dea0d Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 2 Jun 2022 14:05:11 +0200 Subject: [PATCH 05/18] write BrokerRegistrationChangeRecord with InControlledShutdown state when broker requests a controlled shutdown --- .../controller/ClusterControlManager.java | 6 ++ .../controller/ReplicationControlManager.java | 30 +++++++-- .../ReplicationControlManagerTest.java | 64 +++++++++++++++++++ 3 files changed, 96 insertions(+), 4 deletions(-) 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 c2e322781d10a..1be5612f177a1 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -547,6 +547,12 @@ public boolean inControlledShutdown(int brokerId) { return registration.inControlledShutdown(); } + public boolean active(int brokerId) { + BrokerRegistration registration = brokerRegistrations.get(brokerId); + if (registration == null) return false; + return !registration.inControlledShutdown() && !registration.fenced(); + } + BrokerHeartbeatManager heartbeatManager() { if (heartbeatManager == null) { throw new RuntimeException("ClusterControlManager is not active."); diff --git a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java index 13e41c978c6b8..1eb33cce563bf 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java @@ -78,6 +78,7 @@ import org.apache.kafka.metadata.BrokerHeartbeatReply; import org.apache.kafka.metadata.BrokerRegistration; import org.apache.kafka.metadata.BrokerRegistrationFencingChange; +import org.apache.kafka.metadata.BrokerRegistrationInControlledShutdownChange; import org.apache.kafka.metadata.KafkaConfigSchema; import org.apache.kafka.metadata.LeaderRecoveryState; import org.apache.kafka.metadata.PartitionRegistration; @@ -1138,7 +1139,7 @@ void handleBrokerUnregistered(int brokerId, long brokerEpoch, /** * Generate the appropriate records to handle a broker becoming unfenced. * - * First, we create an UnfenceBrokerRecord. Then, we check if if there are any + * First, we create an UnfenceBrokerRecord. Then, we check if there are any * partitions that don't currently have a leader that should be led by the newly * unfenced broker. * @@ -1160,6 +1161,28 @@ void handleBrokerUnfenced(int brokerId, long brokerEpoch, List records) { + if (featureControl.metadataVersion().isInControlledShutdownStateSupported()) { + records.add(new ApiMessageAndVersion(new BrokerRegistrationChangeRecord(). + setBrokerId(brokerId).setBrokerEpoch(brokerEpoch). + setInControlledShutdown(BrokerRegistrationInControlledShutdownChange.IN_CONTROLLED_SHUTDOWN.value()), + (short) 1)); + } + generateLeaderAndIsrUpdates("enterControlledShutdown[" + brokerId + "]", + brokerId, NO_LEADER, records, brokersToIsrs.partitionsWithBrokerInIsr(brokerId)); + } + ControllerResult electLeaders(ElectLeadersRequestData request) { ElectionType electionType = electionType(request.electionType()); List records = new ArrayList<>(); @@ -1278,8 +1301,7 @@ ControllerResult processBrokerHeartbeat( handleBrokerUnfenced(brokerId, brokerEpoch, records); break; case CONTROLLED_SHUTDOWN: - generateLeaderAndIsrUpdates("enterControlledShutdown[" + brokerId + "]", - brokerId, NO_LEADER, records, brokersToIsrs.partitionsWithBrokerInIsr(brokerId)); + handleBrokerInControlledShutdown(brokerId, brokerEpoch, records); break; case SHUTDOWN_NOW: handleBrokerFenced(brokerId, records); @@ -1554,7 +1576,7 @@ void generateLeaderAndIsrUpdates(String context, TopicControlInfo topic = topics.get(topicIdPart.topicId()); if (topic == null) { throw new RuntimeException("Topic ID " + topicIdPart.topicId() + - " existed in isrMembers, but not in the topics map."); + " existed in isrMembers, but not in the topics map."); } PartitionRegistration partition = topic.parts.get(topicIdPart.partitionId()); if (partition == null) { diff --git a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java index d7141a76dd264..bde8f4d8053a8 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java @@ -17,6 +17,7 @@ package org.apache.kafka.controller; +import org.apache.kafka.clients.ApiVersions; import org.apache.kafka.common.ElectionType; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.Uuid; @@ -54,6 +55,7 @@ import org.apache.kafka.common.message.ListPartitionReassignmentsResponseData; import org.apache.kafka.common.message.ListPartitionReassignmentsResponseData.OngoingPartitionReassignment; import org.apache.kafka.common.message.ListPartitionReassignmentsResponseData.OngoingTopicReassignment; +import org.apache.kafka.common.metadata.BrokerRegistrationChangeRecord; import org.apache.kafka.common.metadata.ConfigRecord; import org.apache.kafka.common.metadata.PartitionChangeRecord; import org.apache.kafka.common.metadata.PartitionRecord; @@ -68,6 +70,7 @@ import org.apache.kafka.controller.ReplicationControlManager.KRaftClusterDescriber; import org.apache.kafka.metadata.BrokerHeartbeatReply; import org.apache.kafka.metadata.BrokerRegistration; +import org.apache.kafka.metadata.BrokerRegistrationInControlledShutdownChange; import org.apache.kafka.metadata.LeaderRecoveryState; import org.apache.kafka.metadata.MockRandom; import org.apache.kafka.metadata.PartitionRegistration; @@ -76,11 +79,13 @@ import org.apache.kafka.metadata.placement.StripedReplicaPlacer; import org.apache.kafka.metadata.placement.UsableBroker; import org.apache.kafka.server.common.ApiMessageAndVersion; +import org.apache.kafka.server.common.MetadataVersion; import org.apache.kafka.server.policy.CreateTopicPolicy; import org.apache.kafka.timeline.SnapshotRegistry; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; import org.junit.jupiter.params.provider.ValueSource; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -118,6 +123,7 @@ import static org.apache.kafka.common.protocol.Errors.UNKNOWN_TOPIC_ID; import static org.apache.kafka.common.protocol.Errors.UNKNOWN_TOPIC_OR_PARTITION; import static org.apache.kafka.metadata.LeaderConstants.NO_LEADER; +import static org.apache.kafka.server.common.MetadataVersion.IBP_3_3_IV3; import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; @@ -162,7 +168,23 @@ void replay(List records) throws Exception { this(Optional.empty()); } + ReplicationControlTestContext(MetadataVersion metadataVersion) { + this(metadataVersion, Optional.empty()); + } + ReplicationControlTestContext(Optional createTopicPolicy) { + this(MetadataVersion.latest(), createTopicPolicy); + } + + ReplicationControlTestContext(MetadataVersion metadataVersion, Optional createTopicPolicy) { + FeatureControlManager featureControl = new FeatureControlManager.Builder(). + setSnapshotRegistry(snapshotRegistry). + setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), + QuorumFeatures.defaultFeatureMap(), + Collections.singletonList(0))). + setMetadataVersion(metadataVersion). + build(); + this.replicationControl = new ReplicationControlManager.Builder(). setSnapshotRegistry(snapshotRegistry). setLogContext(logContext). @@ -171,6 +193,7 @@ void replay(List records) throws Exception { setClusterControl(clusterControl). setControllerMetrics(metrics). setCreateTopicPolicy(createTopicPolicy). + setFeatureControl(featureControl). build(); clusterControl.activate(); } @@ -1797,4 +1820,45 @@ public void testKRaftClusterDescriber() throws Exception { new UsableBroker(3, Optional.empty(), false), new UsableBroker(4, Optional.empty(), false))), brokers); } + + @ParameterizedTest + @EnumSource(value = MetadataVersion.class, names = {"IBP_3_3_IV2", "IBP_3_3_IV3"}) + public void testProcessBrokerHeartbeatInControlledShutdown(MetadataVersion metadataVersion) throws Exception { + ReplicationControlTestContext ctx = new ReplicationControlTestContext(metadataVersion); + ctx.registerBrokers(0, 1, 2); + ctx.unfenceBrokers(0, 1, 2); + + Uuid topicId = ctx.createTestTopic("foo", new int[][]{new int[]{0, 1, 2}}).topicId(); + + BrokerHeartbeatRequestData heartbeatRequest = new BrokerHeartbeatRequestData() + .setBrokerId(0) + .setBrokerEpoch(100) + .setCurrentMetadataOffset(0) + .setWantShutDown(true); + + ControllerResult result = ctx.replicationControl + .processBrokerHeartbeat(heartbeatRequest, 0); + + List expectedRecords = new ArrayList<>(); + + if (metadataVersion.isAtLeast(IBP_3_3_IV3)) { + expectedRecords.add(new ApiMessageAndVersion( + new BrokerRegistrationChangeRecord() + .setBrokerEpoch(100) + .setBrokerId(0) + .setInControlledShutdown(BrokerRegistrationInControlledShutdownChange + .IN_CONTROLLED_SHUTDOWN.value()), + (short) 1)); + } + + expectedRecords.add(new ApiMessageAndVersion( + new PartitionChangeRecord() + .setPartitionId(0) + .setTopicId(topicId) + .setIsr(asList(1, 2)) + .setLeader(1), + (short) 0)); + + assertEquals(expectedRecords, result.records()); + } } From e1f2d606c380fe03df9ecdcc8bd19aae49e64510 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 2 Jun 2022 15:32:04 +0200 Subject: [PATCH 06/18] enforce invariants --- .../controller/ReplicationControlManager.java | 55 ++++--- .../ReplicationControlManagerTest.java | 147 ++++++++++++++++-- 2 files changed, 171 insertions(+), 31 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java index 1eb33cce563bf..52bc28e56d019 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java @@ -653,11 +653,11 @@ private ApiError createTopic(CreatableTopic topic, validateManualPartitionAssignment(assignment.brokerIds(), replicationFactor); replicationFactor = OptionalInt.of(assignment.brokerIds().size()); List isr = assignment.brokerIds().stream(). - filter(clusterControl::unfenced).collect(Collectors.toList()); + filter(clusterControl::active).collect(Collectors.toList()); if (isr.isEmpty()) { return new ApiError(Errors.INVALID_REPLICA_ASSIGNMENT, "All brokers specified in the manual partition assignment for " + - "partition " + assignment.partitionIndex() + " are fenced."); + "partition " + assignment.partitionIndex() + " are fenced or in controlled shutdown."); } newParts.put(assignment.partitionIndex(), new PartitionRegistration( Replicas.toArray(assignment.brokerIds()), Replicas.toArray(isr), @@ -683,25 +683,35 @@ private ApiError createTopic(CreatableTopic topic, short replicationFactor = topic.replicationFactor() == -1 ? defaultReplicationFactor : topic.replicationFactor(); try { - List> replicas = clusterControl.replicaPlacer().place(new PlacementSpec( + List> partitions = clusterControl.replicaPlacer().place(new PlacementSpec( 0, numPartitions, replicationFactor ), clusterDescriber); - for (int partitionId = 0; partitionId < replicas.size(); partitionId++) { - int[] r = Replicas.toArray(replicas.get(partitionId)); + for (int partitionId = 0; partitionId < partitions.size(); partitionId++) { + List replicas = partitions.get(partitionId); + List isr = replicas.stream(). + filter(clusterControl::active).collect(Collectors.toList()); + // We need to have at least one replica in the ISR. + if (isr.isEmpty()) isr.add(replicas.get(0)); newParts.put(partitionId, - new PartitionRegistration(r, r, Replicas.NONE, Replicas.NONE, r[0], LeaderRecoveryState.RECOVERED, 0, 0)); + new PartitionRegistration( + Replicas.toArray(replicas), + Replicas.toArray(isr), + Replicas.NONE, + Replicas.NONE, + isr.get(0), + LeaderRecoveryState.RECOVERED, + 0, + 0)); } } catch (InvalidReplicationFactorException e) { return new ApiError(Errors.INVALID_REPLICATION_FACTOR, "Unable to replicate the partition " + replicationFactor + " time(s): " + e.getMessage()); } - ApiError error = maybeCheckCreateTopicPolicy(() -> { - return new CreateTopicPolicy.RequestMetadata( - topic.name(), numPartitions, replicationFactor, null, creationConfigs); - }); + ApiError error = maybeCheckCreateTopicPolicy(() -> new CreateTopicPolicy.RequestMetadata( + topic.name(), numPartitions, replicationFactor, null, creationConfigs)); if (error.isFailure()) return error; } Uuid topicId = Uuid.randomUuid(); @@ -938,7 +948,7 @@ ControllerResult alterPartition(AlterPartitionReques partition, topic.id, partitionId, - r -> clusterControl.unfenced(r), + clusterControl::active, featureControl.metadataVersion().isLeaderRecoverySupported()); if (configurationControl.uncleanLeaderElectionEnabledForTopic(topicData.name())) { builder.setElection(PartitionChangeBuilder.Election.UNCLEAN); @@ -1268,7 +1278,7 @@ ApiError electLeader(String topic, int partitionId, ElectionType electionType, PartitionChangeBuilder builder = new PartitionChangeBuilder(partition, topicId, partitionId, - r -> clusterControl.unfenced(r), + clusterControl::active, featureControl.metadataVersion().isLeaderRecoverySupported()); builder.setElection(election); Optional record = builder.build(); @@ -1383,7 +1393,7 @@ ControllerResult maybeBalancePartitionLeaders() { partition, topicPartition.topicId(), topicPartition.partitionId(), - r -> clusterControl.unfenced(r), + clusterControl::active, featureControl.metadataVersion().isLeaderRecoverySupported() ); builder.setElection(PartitionChangeBuilder.Election.PREFERRED); @@ -1471,11 +1481,11 @@ void createPartitions(CreatePartitionsTopic topic, OptionalInt.of(replicationFactor)); placements.add(assignment.brokerIds()); List isr = assignment.brokerIds().stream(). - filter(clusterControl::unfenced).collect(Collectors.toList()); + filter(clusterControl::active).collect(Collectors.toList()); if (isr.isEmpty()) { throw new InvalidReplicaAssignmentException( "All brokers specified in the manual partition assignment for " + - "partition " + (startPartitionId + i) + " are fenced."); + "partition " + (startPartitionId + i) + " are fenced or in controlled shutdown."); } isrs.add(isr); } @@ -1489,12 +1499,15 @@ void createPartitions(CreatePartitionsTopic topic, } int partitionId = startPartitionId; for (int i = 0; i < placements.size(); i++) { - List placement = placements.get(i); - List isr = isrs.get(i); + List replicas = placements.get(i); + List isr = isrs.get(i).stream(). + filter(clusterControl::active).collect(Collectors.toList()); + // We need to have at least one replica in the ISR. + if (isr.isEmpty()) isr.add(replicas.get(0)); records.add(new ApiMessageAndVersion(new PartitionRecord(). setPartitionId(partitionId). setTopicId(topicId). - setReplicas(placement). + setReplicas(replicas). setIsr(isr). setLeaderRecoveryState(LeaderRecoveryState.RECOVERED.value()). setRemovingReplicas(Collections.emptyList()). @@ -1569,7 +1582,7 @@ void generateLeaderAndIsrUpdates(String context, // where there is an unclean leader election which chooses a leader from outside // the ISR. Function isAcceptableLeader = - r -> (r != brokerToRemove) && (r == brokerToAdd || clusterControl.unfenced(r)); + r -> (r != brokerToRemove) && (r == brokerToAdd || clusterControl.active(r)); while (iterator.hasNext()) { TopicIdPartition topicIdPart = iterator.next(); @@ -1696,7 +1709,7 @@ Optional cancelPartitionReassignment(String topicName, PartitionChangeBuilder builder = new PartitionChangeBuilder(part, tp.topicId(), tp.partitionId(), - r -> clusterControl.unfenced(r), + clusterControl::active, featureControl.metadataVersion().isLeaderRecoverySupported()); if (configurationControl.uncleanLeaderElectionEnabledForTopic(topicName)) { builder.setElection(PartitionChangeBuilder.Election.UNCLEAN); @@ -1748,7 +1761,7 @@ Optional changePartitionReassignment(TopicIdPartition tp, PartitionChangeBuilder builder = new PartitionChangeBuilder(part, tp.topicId(), tp.partitionId(), - r -> clusterControl.unfenced(r), + clusterControl::active, featureControl.metadataVersion().isLeaderRecoverySupported()); if (!reassignment.merged().equals(currentReplicas)) { builder.setTargetReplicas(reassignment.merged()); diff --git a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java index bde8f4d8053a8..1dee7bfad5e70 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java @@ -322,7 +322,7 @@ void alterPartition( replay(alterPartition.records()); } - void unfenceBrokers(Integer... brokerIds) throws Exception { + void unfenceBrokers(Integer... brokerIds) throws Exception { unfenceBrokers(Utils.mkSet(brokerIds)); } @@ -339,6 +339,20 @@ void unfenceBrokers(Set brokerIds) throws Exception { } } + void inControlledShutdownBrokers(Integer... brokerIds) throws Exception { + inControlledShutdownBrokers(Utils.mkSet(brokerIds)); + } + + void inControlledShutdownBrokers(Set brokerIds) throws Exception { + for (int brokerId : brokerIds) { + BrokerRegistrationChangeRecord record = new BrokerRegistrationChangeRecord() + .setBrokerId(brokerId) + .setBrokerEpoch(brokerId + 100) + .setInControlledShutdown(BrokerRegistrationInControlledShutdownChange.IN_CONTROLLED_SHUTDOWN.value()); + replay(singletonList(new ApiMessageAndVersion(record, (short) 1))); + } + } + void alterTopicConfig( String topic, String configKey, @@ -467,6 +481,46 @@ public void testCreateTopics() throws Exception { ctx.replicationControl.iterator(Long.MAX_VALUE)); } + @Test + public void testCreateTopicsInvariants() throws Exception { + ReplicationControlTestContext ctx = new ReplicationControlTestContext(); + ReplicationControlManager replicationControl = ctx.replicationControl; + + CreateTopicsRequestData request = new CreateTopicsRequestData(); + request.topics().add(new CreatableTopic().setName("foo"). + setNumPartitions(-1).setReplicationFactor((short) -1)); + + ctx.registerBrokers(0, 1, 2); + ctx.unfenceBrokers(0, 1); + ctx.inControlledShutdownBrokers(1); + + ControllerResult result = + replicationControl.createTopics(request, Collections.singleton("foo")); + + CreateTopicsResponseData expectedResponse = new CreateTopicsResponseData(); + expectedResponse.topics().add(new CreatableTopicResult().setName("foo"). + setNumPartitions(1).setReplicationFactor((short) 3). + setErrorMessage(null).setErrorCode((short) 0). + setTopicId(result.response().topics().find("foo").topicId())); + assertEquals(expectedResponse, result.response()); + + ctx.replay(result.records()); + + // Broker 2 cannot be in the ISR because it is fenced and broker 1 + // cannot be in the ISR because it is in controlled shutdown. + assertEquals( + new PartitionRegistration(new int[]{1, 0, 2}, + new int[]{0}, + Replicas.NONE, + Replicas.NONE, + 0, + LeaderRecoveryState.RECOVERED, + 0, + 0), + replicationControl.getPartition( + ((TopicRecord) result.records().get(0).message()).topicId(), 0)); + } + @Test public void testCreateTopicsWithConfigs() throws Exception { ReplicationControlTestContext ctx = new ReplicationControlTestContext(); @@ -1123,6 +1177,45 @@ public void testCreatePartitions() throws Exception { ctx.replay(createPartitionsResult2.records()); } + @Test + public void testCreatePartitionsInvariants() throws Exception { + ReplicationControlTestContext ctx = new ReplicationControlTestContext(); + ReplicationControlManager replicationControl = ctx.replicationControl; + + CreateTopicsRequestData request = new CreateTopicsRequestData(); + request.topics().add(new CreatableTopic().setName("foo"). + setNumPartitions(1).setReplicationFactor((short) 3)); + + ctx.registerBrokers(0, 1, 2); + ctx.unfenceBrokers(0, 1); + ctx.inControlledShutdownBrokers(1); + + ControllerResult result = + replicationControl.createTopics(request, Collections.singleton("foo")); + ctx.replay(result.records()); + + List topics = asList(new CreatePartitionsTopic(). + setName("foo").setCount(2).setAssignments(null)); + + ControllerResult> createPartitionsResult = + replicationControl.createPartitions(topics); + ctx.replay(createPartitionsResult.records()); + + // Broker 2 cannot be in the ISR because it is fenced and broker 1 + // cannot be in the ISR because it is in controlled shutdown. + assertEquals( + new PartitionRegistration(new int[]{0, 1, 2}, + new int[]{0}, + Replicas.NONE, + Replicas.NONE, + 0, + LeaderRecoveryState.RECOVERED, + 0, + 0), + replicationControl.getPartition( + ((TopicRecord) result.records().get(0).message()).topicId(), 1)); + } + @Test public void testValidateGoodManualPartitionAssignments() throws Exception { ReplicationControlTestContext ctx = new ReplicationControlTestContext(); @@ -1593,14 +1686,15 @@ public void testElectPreferredLeaders() throws Exception { ReplicationControlTestContext ctx = new ReplicationControlTestContext(); ReplicationControlManager replication = ctx.replicationControl; ctx.registerBrokers(0, 1, 2, 3, 4); - ctx.unfenceBrokers(2, 3, 4); + ctx.unfenceBrokers(1, 2, 3, 4); + ctx.inControlledShutdownBrokers(1); Uuid fooId = ctx.createTestTopic("foo", new int[][]{ new int[]{1, 2, 3}, new int[]{2, 3, 4}, new int[]{0, 2, 1}}).topicId(); ElectLeadersRequestData request1 = new ElectLeadersRequestData(). setElectionType(ElectionType.PREFERRED.value). setTopicPartitions(new TopicPartitionsCollection(asList( new TopicPartitions().setTopic("foo"). - setPartitions(asList(0, 1)), + setPartitions(asList(0, 1, 2)), new TopicPartitions().setTopic("bar"). setPartitions(asList(0, 1))).iterator())); ControllerResult election1Result = @@ -1614,6 +1708,10 @@ public void testElectPreferredLeaders() throws Exception { new TopicPartition("foo", 1), new ApiError(ELECTION_NOT_NEEDED) ), + Utils.mkEntry( + new TopicPartition("foo", 2), + new ApiError(PREFERRED_LEADER_NOT_AVAILABLE) + ), Utils.mkEntry( new TopicPartition("bar", 0), new ApiError(UNKNOWN_TOPIC_OR_PARTITION, "No such topic as bar") @@ -1625,14 +1723,21 @@ public void testElectPreferredLeaders() throws Exception { )); assertElectLeadersResponse(expectedResponse1, election1Result.response()); assertEquals(Collections.emptyList(), election1Result.records()); + + // Broker 1 must be registered to get out from the controlled shutdown state. + ctx.registerBrokers(1); ctx.unfenceBrokers(0, 1); ControllerResult alterPartitionResult = replication.alterPartition( new AlterPartitionRequestData().setBrokerId(2).setBrokerEpoch(102). setTopics(asList(new AlterPartitionRequestData.TopicData().setName("foo"). - setPartitions(asList(new AlterPartitionRequestData.PartitionData(). - setPartitionIndex(0).setPartitionEpoch(0). - setLeaderEpoch(0).setNewIsr(asList(1, 2, 3))))))); + setPartitions(asList( + new AlterPartitionRequestData.PartitionData(). + setPartitionIndex(0).setPartitionEpoch(0). + setLeaderEpoch(0).setNewIsr(asList(1, 2, 3)), + new AlterPartitionRequestData.PartitionData(). + setPartitionIndex(2).setPartitionEpoch(0). + setLeaderEpoch(0).setNewIsr(asList(0, 2, 1))))))); assertEquals(new AlterPartitionResponseData().setTopics(asList( new AlterPartitionResponseData.TopicData().setName("foo").setPartitions(asList( new AlterPartitionResponseData.PartitionData(). @@ -1641,6 +1746,13 @@ public void testElectPreferredLeaders() throws Exception { setLeaderEpoch(0). setIsr(asList(1, 2, 3)). setPartitionEpoch(1). + setErrorCode(NONE.code()), + new AlterPartitionResponseData.PartitionData(). + setPartitionIndex(2). + setLeaderId(2). + setLeaderEpoch(0). + setIsr(asList(0, 2, 1)). + setPartitionEpoch(1). setErrorCode(NONE.code()))))), alterPartitionResult.response()); @@ -1653,6 +1765,10 @@ public void testElectPreferredLeaders() throws Exception { new TopicPartition("foo", 1), new ApiError(ELECTION_NOT_NEEDED) ), + Utils.mkEntry( + new TopicPartition("foo", 2), + ApiError.NONE + ), Utils.mkEntry( new TopicPartition("bar", 0), new ApiError(UNKNOWN_TOPIC_OR_PARTITION, "No such topic as bar") @@ -1667,10 +1783,21 @@ public void testElectPreferredLeaders() throws Exception { ControllerResult election2Result = replication.electLeaders(request1); assertElectLeadersResponse(expectedResponse2, election2Result.response()); - assertEquals(asList(new ApiMessageAndVersion(new PartitionChangeRecord(). - setPartitionId(0). - setTopicId(fooId). - setLeader(1), (short) 0)), election2Result.records()); + assertEquals( + asList( + new ApiMessageAndVersion( + new PartitionChangeRecord(). + setPartitionId(0). + setTopicId(fooId). + setLeader(1), + (short) 0), + new ApiMessageAndVersion( + new PartitionChangeRecord(). + setPartitionId(2). + setTopicId(fooId). + setLeader(0), + (short) 0)), + election2Result.records()); } @Test From bbc95daf86d1e332c54362b0e1536553ed5c1c56 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 2 Jun 2022 16:26:49 +0200 Subject: [PATCH 07/18] fix core tests --- .../kafka/server/metadata/BrokerMetadataListenerTest.scala | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/core/src/test/scala/unit/kafka/server/metadata/BrokerMetadataListenerTest.scala b/core/src/test/scala/unit/kafka/server/metadata/BrokerMetadataListenerTest.scala index 20df47a6a6ab0..948e05133701b 100644 --- a/core/src/test/scala/unit/kafka/server/metadata/BrokerMetadataListenerTest.scala +++ b/core/src/test/scala/unit/kafka/server/metadata/BrokerMetadataListenerTest.scala @@ -98,11 +98,11 @@ class BrokerMetadataListenerTest { assertEquals(200L, newImage.highestOffsetAndEpoch().offset) assertEquals(new BrokerRegistration(0, 100L, Uuid.fromString("GFBwlTcpQUuLYQ2ig05CSg"), Collections.emptyList[Endpoint](), - Collections.emptyMap[String, VersionRange](), Optional.empty[String](), false), + Collections.emptyMap[String, VersionRange](), Optional.empty[String](), false, false), delta.clusterDelta().broker(0)) assertEquals(new BrokerRegistration(1, 200L, Uuid.fromString("QkOQtNKVTYatADcaJ28xDg"), Collections.emptyList[Endpoint](), - Collections.emptyMap[String, VersionRange](), Optional.empty[String](), true), + Collections.emptyMap[String, VersionRange](), Optional.empty[String](), true, false), delta.clusterDelta().broker(1)) } From 3ea21262f57201c288eca45dff9a0d683950778b Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 2 Jun 2022 16:59:17 +0200 Subject: [PATCH 08/18] refactor + image test --- .../controller/ClusterControlManager.java | 9 +++-- .../org/apache/kafka/image/ClusterDelta.java | 34 ++++++++++++++++--- .../kafka/metadata/BrokerRegistration.java | 17 ++++++++-- .../apache/kafka/image/ClusterImageTest.java | 8 ++++- 4 files changed, 55 insertions(+), 13 deletions(-) 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 1be5612f177a1..04b7d3ba03c08 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -476,11 +476,10 @@ private void replayRegistrationChange( "registration with that epoch found", record.toString())); } else { BrokerRegistration nextRegistration = curRegistration; - if (fencingChange != BrokerRegistrationFencingChange.NONE) { - nextRegistration = nextRegistration.cloneWithFencing(fencingChange.asBoolean().get()); - } - if (inControlledShutdownChange != BrokerRegistrationInControlledShutdownChange.NONE) { - nextRegistration = nextRegistration.cloneWithInControlledShutdown(inControlledShutdownChange.asBoolean().get()); + if (fencingChange != BrokerRegistrationFencingChange.NONE + || inControlledShutdownChange != BrokerRegistrationInControlledShutdownChange.NONE) { + nextRegistration = nextRegistration.cloneWith( + fencingChange.asBoolean(), inControlledShutdownChange.asBoolean()); } if (!curRegistration.equals(nextRegistration)) { brokerRegistrations.put(brokerId, nextRegistration); 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 1c4d66b9e922e..14bb8620ac4fe 100644 --- a/metadata/src/main/java/org/apache/kafka/image/ClusterDelta.java +++ b/metadata/src/main/java/org/apache/kafka/image/ClusterDelta.java @@ -23,6 +23,8 @@ import org.apache.kafka.common.metadata.UnfenceBrokerRecord; import org.apache.kafka.common.metadata.UnregisterBrokerRecord; import org.apache.kafka.metadata.BrokerRegistration; +import org.apache.kafka.metadata.BrokerRegistrationFencingChange; +import org.apache.kafka.metadata.BrokerRegistrationInControlledShutdownChange; import org.apache.kafka.server.common.MetadataVersion; import java.util.HashMap; @@ -102,10 +104,34 @@ public void replay(UnfenceBrokerRecord record) { public void replay(BrokerRegistrationChangeRecord record) { BrokerRegistration broker = getBrokerOrThrow(record.brokerId(), record.brokerEpoch(), "change"); - if (record.fenced() < 0) { - changedBrokers.put(record.brokerId(), Optional.of(broker.cloneWithFencing(false))); - } else if (record.fenced() > 0) { - changedBrokers.put(record.brokerId(), Optional.of(broker.cloneWithFencing(true))); + Optional fencingChange = + BrokerRegistrationFencingChange.fromValue(record.fenced()); + if (!fencingChange.isPresent()) { + throw new IllegalStateException(String.format("Unable to replay %s: unknown " + + "value for fenced field: %d", record.toString(), record.fenced())); + } + Optional inControlledShutdownChange = + BrokerRegistrationInControlledShutdownChange.fromValue(record.inControlledShutdown()); + if (!inControlledShutdownChange.isPresent()) { + throw new IllegalStateException(String.format("Unable to replay %s: unknown " + + "value for inControlledShutdown field: %d", record.toString(), record.inControlledShutdown())); + } + + replayRegistration(record.brokerId(), broker, fencingChange.get(), inControlledShutdownChange.get()); + } + + private void replayRegistration( + int brokerId, + BrokerRegistration broker, + BrokerRegistrationFencingChange fencingChange, + BrokerRegistrationInControlledShutdownChange inControlledShutdownChange + ) { + if (fencingChange != BrokerRegistrationFencingChange.NONE + || inControlledShutdownChange != BrokerRegistrationInControlledShutdownChange.NONE) { + changedBrokers.put(brokerId, Optional.of(broker.cloneWith( + fencingChange.asBoolean(), + inControlledShutdownChange.asBoolean() + ))); } } diff --git a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java index 64d967c68b378..2a6b39220508b 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java @@ -231,8 +231,19 @@ public BrokerRegistration cloneWithFencing(boolean fencing) { supportedFeatures, rack, fencing, inControlledShutdown); } - public BrokerRegistration cloneWithInControlledShutdown(boolean inControlledShutdown) { - return new BrokerRegistration(id, epoch, incarnationId, listeners, - supportedFeatures, rack, fenced, inControlledShutdown); + public BrokerRegistration cloneWith( + Optional fencingChange, + Optional inControlledShutdownChange + ) { + return new BrokerRegistration( + id, + epoch, + incarnationId, + listeners, + supportedFeatures, + rack, + fencingChange.orElse(fenced), + inControlledShutdownChange.orElse(inControlledShutdown) + ); } } diff --git a/metadata/src/test/java/org/apache/kafka/image/ClusterImageTest.java b/metadata/src/test/java/org/apache/kafka/image/ClusterImageTest.java index d1356f66b64b8..c1a1886b3da87 100644 --- a/metadata/src/test/java/org/apache/kafka/image/ClusterImageTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/ClusterImageTest.java @@ -19,11 +19,13 @@ import org.apache.kafka.common.Endpoint; import org.apache.kafka.common.Uuid; +import org.apache.kafka.common.metadata.BrokerRegistrationChangeRecord; import org.apache.kafka.common.metadata.FenceBrokerRecord; import org.apache.kafka.common.metadata.UnfenceBrokerRecord; import org.apache.kafka.common.metadata.UnregisterBrokerRecord; import org.apache.kafka.common.security.auth.SecurityProtocol; import org.apache.kafka.metadata.BrokerRegistration; +import org.apache.kafka.metadata.BrokerRegistrationInControlledShutdownChange; import org.apache.kafka.metadata.RecordTestUtils; import org.apache.kafka.metadata.VersionRange; import org.apache.kafka.server.common.ApiMessageAndVersion; @@ -87,6 +89,10 @@ public class ClusterImageTest { setId(0).setEpoch(1000), UNFENCE_BROKER_RECORD.highestSupportedVersion())); DELTA1_RECORDS.add(new ApiMessageAndVersion(new FenceBrokerRecord(). setId(1).setEpoch(1001), FENCE_BROKER_RECORD.highestSupportedVersion())); + DELTA1_RECORDS.add(new ApiMessageAndVersion(new BrokerRegistrationChangeRecord(). + setBrokerId(0).setBrokerEpoch(1000).setInControlledShutdown( + BrokerRegistrationInControlledShutdownChange.IN_CONTROLLED_SHUTDOWN.value()), + FENCE_BROKER_RECORD.highestSupportedVersion())); DELTA1_RECORDS.add(new ApiMessageAndVersion(new UnregisterBrokerRecord(). setBrokerId(2).setBrokerEpoch(123), UNREGISTER_BROKER_RECORD.highestSupportedVersion())); @@ -102,7 +108,7 @@ public class ClusterImageTest { Collections.singletonMap("foo", VersionRange.of((short) 1, (short) 3)), Optional.empty(), false, - false)); + true)); map2.put(1, new BrokerRegistration(1, 1001, Uuid.fromString("U52uRe20RsGI0RvpcTx33Q"), From d9a1af82a85e9422413d96445c8165367709f959 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 2 Jun 2022 20:58:23 +0200 Subject: [PATCH 09/18] cleanups and fixes --- .../metadata/BrokerMetadataListener.scala | 2 +- .../metadata/BrokerMetadataSnapshotter.scala | 5 ++--- .../controller/ClusterControlManager.java | 14 ++++++++------ .../org/apache/kafka/image/ClusterImage.java | 6 +++--- .../org/apache/kafka/image/MetadataImage.java | 5 +++-- .../kafka/metadata/BrokerRegistration.java | 19 ++++++++++++++----- .../apache/kafka/image/ClusterImageTest.java | 3 ++- .../apache/kafka/image/MetadataImageTest.java | 3 ++- .../metadata/BrokerRegistrationTest.java | 5 +++-- 9 files changed, 38 insertions(+), 24 deletions(-) diff --git a/core/src/main/scala/kafka/server/metadata/BrokerMetadataListener.scala b/core/src/main/scala/kafka/server/metadata/BrokerMetadataListener.scala index fa0bc52d7aa01..64fd23bd971d8 100644 --- a/core/src/main/scala/kafka/server/metadata/BrokerMetadataListener.scala +++ b/core/src/main/scala/kafka/server/metadata/BrokerMetadataListener.scala @@ -322,7 +322,7 @@ class BrokerMetadataListener( } override def run(): Unit = { - _image.write(this) + _image.write(this, _image.features().metadataVersion()) future.complete(records) } } diff --git a/core/src/main/scala/kafka/server/metadata/BrokerMetadataSnapshotter.scala b/core/src/main/scala/kafka/server/metadata/BrokerMetadataSnapshotter.scala index f32c4d3238f16..e2a69b159255a 100644 --- a/core/src/main/scala/kafka/server/metadata/BrokerMetadataSnapshotter.scala +++ b/core/src/main/scala/kafka/server/metadata/BrokerMetadataSnapshotter.scala @@ -17,7 +17,6 @@ package kafka.server.metadata import java.util.concurrent.RejectedExecutionException - import kafka.utils.Logging import org.apache.kafka.image.MetadataImage import org.apache.kafka.common.utils.{LogContext, Time} @@ -25,7 +24,6 @@ import org.apache.kafka.queue.{EventQueue, KafkaEventQueue} import org.apache.kafka.server.common.ApiMessageAndVersion import org.apache.kafka.snapshot.SnapshotWriter - trait SnapshotWriterBuilder { def build(committedOffset: Long, committedEpoch: Int, @@ -75,7 +73,8 @@ class BrokerMetadataSnapshotter( extends EventQueue.Event { override def run(): Unit = { try { - image.write(writer.append(_)) + val metadataVersion = image.features().metadataVersion(); + image.write(writer.append(_), metadataVersion) writer.freeze() } finally { try { 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 04b7d3ba03c08..b86ad8440a081 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -54,7 +54,6 @@ import org.slf4j.Logger; import java.util.ArrayList; -import java.util.Collections; import java.util.HashMap; import java.util.Iterator; import java.util.List; @@ -68,6 +67,7 @@ import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; +import static java.util.Collections.singletonList; import static java.util.concurrent.TimeUnit.NANOSECONDS; @@ -151,7 +151,7 @@ ClusterControlManager build() { setSnapshotRegistry(snapshotRegistry). setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), QuorumFeatures.defaultFeatureMap(), - Collections.singletonList(0))). + singletonList(0))). setMetadataVersion(MetadataVersion.latest()). build(); } @@ -610,16 +610,18 @@ public List next() { setMaxSupportedVersion(featureEntry.getValue().max()). setMinSupportedVersion(featureEntry.getValue().min())); } - List batch = new ArrayList<>(); - batch.add(new ApiMessageAndVersion(new RegisterBrokerRecord(). + RegisterBrokerRecord record = new RegisterBrokerRecord(). setBrokerId(brokerId). setIncarnationId(registration.incarnationId()). setBrokerEpoch(registration.epoch()). setEndPoints(endpoints). setFeatures(features). setRack(registration.rack().orElse(null)). - setFenced(registration.fenced()), registerBrokerRecordVersion())); - return batch; + setFenced(registration.fenced()); + if (featureControl.metadataVersion().isInControlledShutdownStateSupported()) { + record.setInControlledShutdown(registration.inControlledShutdown()); + } + return singletonList(new ApiMessageAndVersion(record, registerBrokerRecordVersion())); } } diff --git a/metadata/src/main/java/org/apache/kafka/image/ClusterImage.java b/metadata/src/main/java/org/apache/kafka/image/ClusterImage.java index 3cf36fa0885cd..d513cbca35f88 100644 --- a/metadata/src/main/java/org/apache/kafka/image/ClusterImage.java +++ b/metadata/src/main/java/org/apache/kafka/image/ClusterImage.java @@ -19,6 +19,7 @@ import org.apache.kafka.metadata.BrokerRegistration; import org.apache.kafka.server.common.ApiMessageAndVersion; +import org.apache.kafka.server.common.MetadataVersion; import java.util.ArrayList; import java.util.Collections; @@ -27,7 +28,6 @@ import java.util.function.Consumer; import java.util.stream.Collectors; - /** * Represents the cluster in the metadata image. * @@ -54,10 +54,10 @@ public BrokerRegistration broker(int nodeId) { return brokers.get(nodeId); } - public void write(Consumer> out) { + public void write(Consumer> out, MetadataVersion metadataVersion) { List batch = new ArrayList<>(); for (BrokerRegistration broker : brokers.values()) { - batch.add(broker.toRecord()); + batch.add(broker.toRecord(metadataVersion)); } out.accept(batch); } diff --git a/metadata/src/main/java/org/apache/kafka/image/MetadataImage.java b/metadata/src/main/java/org/apache/kafka/image/MetadataImage.java index 48bed5f8a9b62..1d68374e0612f 100644 --- a/metadata/src/main/java/org/apache/kafka/image/MetadataImage.java +++ b/metadata/src/main/java/org/apache/kafka/image/MetadataImage.java @@ -23,6 +23,7 @@ import java.util.List; import java.util.Objects; import java.util.function.Consumer; +import org.apache.kafka.server.common.MetadataVersion; /** @@ -119,11 +120,11 @@ public AclsImage acls() { return acls; } - public void write(Consumer> out) { + public void write(Consumer> out, MetadataVersion metadataVersion) { // Features should be written out first so we can include the metadata.version at the beginning of the // snapshot features.write(out); - cluster.write(out); + cluster.write(out, metadataVersion); topics.write(out); configs.write(out); clientQuotas.write(out); diff --git a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java index 2a6b39220508b..2c0a9aab67174 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java @@ -20,6 +20,7 @@ import org.apache.kafka.common.Endpoint; import org.apache.kafka.common.Node; import org.apache.kafka.common.Uuid; +import org.apache.kafka.server.common.MetadataVersion; import org.apache.kafka.common.metadata.RegisterBrokerRecord; import org.apache.kafka.common.metadata.RegisterBrokerRecord.BrokerEndpoint; import org.apache.kafka.common.metadata.RegisterBrokerRecord.BrokerFeature; @@ -159,14 +160,18 @@ public boolean inControlledShutdown() { return inControlledShutdown; } - public ApiMessageAndVersion toRecord() { + public ApiMessageAndVersion toRecord(MetadataVersion metadataVersion) { RegisterBrokerRecord registrationRecord = new RegisterBrokerRecord(). setBrokerId(id). setRack(rack.orElse(null)). setBrokerEpoch(epoch). setIncarnationId(incarnationId). - setFenced(fenced). - setInControlledShutdown(inControlledShutdown); + setFenced(fenced); + + if (metadataVersion.isInControlledShutdownStateSupported()) { + registrationRecord.setInControlledShutdown(inControlledShutdown); + } + for (Entry entry : listeners.entrySet()) { Endpoint endpoint = entry.getValue(); registrationRecord.endPoints().add(new BrokerEndpoint(). @@ -175,19 +180,23 @@ public ApiMessageAndVersion toRecord() { setPort(endpoint.port()). setSecurityProtocol(endpoint.securityProtocol().id)); } + for (Entry entry : supportedFeatures.entrySet()) { registrationRecord.features().add(new BrokerFeature(). setName(entry.getKey()). setMinSupportedVersion(entry.getValue().min()). setMaxSupportedVersion(entry.getValue().max())); } - return new ApiMessageAndVersion(registrationRecord, (short) 0); + + short version = metadataVersion.isInControlledShutdownStateSupported() ? + (short) 1 : (short) 0; + return new ApiMessageAndVersion(registrationRecord, version); } @Override public int hashCode() { return Objects.hash(id, epoch, incarnationId, listeners, supportedFeatures, - rack, fenced); + rack, fenced, inControlledShutdown); } @Override diff --git a/metadata/src/test/java/org/apache/kafka/image/ClusterImageTest.java b/metadata/src/test/java/org/apache/kafka/image/ClusterImageTest.java index c1a1886b3da87..59d5d2fed940a 100644 --- a/metadata/src/test/java/org/apache/kafka/image/ClusterImageTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/ClusterImageTest.java @@ -29,6 +29,7 @@ import org.apache.kafka.metadata.RecordTestUtils; import org.apache.kafka.metadata.VersionRange; import org.apache.kafka.server.common.ApiMessageAndVersion; +import org.apache.kafka.server.common.MetadataVersion; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; @@ -142,7 +143,7 @@ public void testImage2RoundTrip() throws Throwable { private void testToImageAndBack(ClusterImage image) throws Throwable { MockSnapshotConsumer writer = new MockSnapshotConsumer(); - image.write(writer); + image.write(writer, MetadataVersion.latest()); ClusterDelta delta = new ClusterDelta(ClusterImage.EMPTY); RecordTestUtils.replayAllBatches(delta, writer.batches()); ClusterImage nextImage = delta.apply(); diff --git a/metadata/src/test/java/org/apache/kafka/image/MetadataImageTest.java b/metadata/src/test/java/org/apache/kafka/image/MetadataImageTest.java index 00193dc7414df..83a395ec11c7a 100644 --- a/metadata/src/test/java/org/apache/kafka/image/MetadataImageTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/MetadataImageTest.java @@ -19,6 +19,7 @@ import org.apache.kafka.metadata.RecordTestUtils; import org.apache.kafka.raft.OffsetAndEpoch; +import org.apache.kafka.server.common.MetadataVersion; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; @@ -86,7 +87,7 @@ public void testImage2RoundTrip() throws Throwable { private void testToImageAndBack(MetadataImage image) throws Throwable { MockSnapshotConsumer writer = new MockSnapshotConsumer(); - image.write(writer); + image.write(writer, MetadataVersion.latest()); MetadataDelta delta = new MetadataDelta(MetadataImage.EMPTY); RecordTestUtils.replayAllBatches( delta, image.highestOffsetAndEpoch().offset, image.highestOffsetAndEpoch().epoch, writer.batches()); diff --git a/metadata/src/test/java/org/apache/kafka/metadata/BrokerRegistrationTest.java b/metadata/src/test/java/org/apache/kafka/metadata/BrokerRegistrationTest.java index 1178cb31ace0b..10d1169412cd3 100644 --- a/metadata/src/test/java/org/apache/kafka/metadata/BrokerRegistrationTest.java +++ b/metadata/src/test/java/org/apache/kafka/metadata/BrokerRegistrationTest.java @@ -23,6 +23,7 @@ import org.apache.kafka.common.metadata.RegisterBrokerRecord; import org.apache.kafka.common.security.auth.SecurityProtocol; import org.apache.kafka.server.common.ApiMessageAndVersion; +import org.apache.kafka.server.common.MetadataVersion; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; @@ -86,11 +87,11 @@ public void testFromRecordAndToRecord() { } private void testRoundTrip(BrokerRegistration registration) { - ApiMessageAndVersion messageAndVersion = registration.toRecord(); + ApiMessageAndVersion messageAndVersion = registration.toRecord(MetadataVersion.latest()); BrokerRegistration registration2 = BrokerRegistration.fromRecord( (RegisterBrokerRecord) messageAndVersion.message()); assertEquals(registration, registration2); - ApiMessageAndVersion messageAndVersion2 = registration2.toRecord(); + ApiMessageAndVersion messageAndVersion2 = registration2.toRecord(MetadataVersion.latest()); assertEquals(messageAndVersion, messageAndVersion2); } From b65e57f3cb072f40c5a8e8ce488c1b6650718777 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Fri, 3 Jun 2022 08:46:15 +0200 Subject: [PATCH 10/18] refactor --- .../kafka/server/metadata/BrokerMetadataListener.scala | 2 +- .../kafka/server/metadata/BrokerMetadataSnapshotter.scala | 3 +-- .../main/java/org/apache/kafka/image/MetadataImage.java | 8 +++++++- .../java/org/apache/kafka/image/MetadataImageTest.java | 3 +-- 4 files changed, 10 insertions(+), 6 deletions(-) diff --git a/core/src/main/scala/kafka/server/metadata/BrokerMetadataListener.scala b/core/src/main/scala/kafka/server/metadata/BrokerMetadataListener.scala index 64fd23bd971d8..fa0bc52d7aa01 100644 --- a/core/src/main/scala/kafka/server/metadata/BrokerMetadataListener.scala +++ b/core/src/main/scala/kafka/server/metadata/BrokerMetadataListener.scala @@ -322,7 +322,7 @@ class BrokerMetadataListener( } override def run(): Unit = { - _image.write(this, _image.features().metadataVersion()) + _image.write(this) future.complete(records) } } diff --git a/core/src/main/scala/kafka/server/metadata/BrokerMetadataSnapshotter.scala b/core/src/main/scala/kafka/server/metadata/BrokerMetadataSnapshotter.scala index e2a69b159255a..b5179c32f1416 100644 --- a/core/src/main/scala/kafka/server/metadata/BrokerMetadataSnapshotter.scala +++ b/core/src/main/scala/kafka/server/metadata/BrokerMetadataSnapshotter.scala @@ -73,8 +73,7 @@ class BrokerMetadataSnapshotter( extends EventQueue.Event { override def run(): Unit = { try { - val metadataVersion = image.features().metadataVersion(); - image.write(writer.append(_), metadataVersion) + image.write(writer.append(_)) writer.freeze() } finally { try { diff --git a/metadata/src/main/java/org/apache/kafka/image/MetadataImage.java b/metadata/src/main/java/org/apache/kafka/image/MetadataImage.java index 1d68374e0612f..d1aa3c6e8a938 100644 --- a/metadata/src/main/java/org/apache/kafka/image/MetadataImage.java +++ b/metadata/src/main/java/org/apache/kafka/image/MetadataImage.java @@ -120,7 +120,13 @@ public AclsImage acls() { return acls; } - public void write(Consumer> out, MetadataVersion metadataVersion) { + public void write(Consumer> out) { + // We use the latest metadata version if this image does not have + // a specific version set. + MetadataVersion metadataVersion = features.metadataVersion(); + if (metadataVersion.equals(MetadataVersion.UNINITIALIZED)) { + metadataVersion = MetadataVersion.latest(); + } // Features should be written out first so we can include the metadata.version at the beginning of the // snapshot features.write(out); diff --git a/metadata/src/test/java/org/apache/kafka/image/MetadataImageTest.java b/metadata/src/test/java/org/apache/kafka/image/MetadataImageTest.java index 83a395ec11c7a..00193dc7414df 100644 --- a/metadata/src/test/java/org/apache/kafka/image/MetadataImageTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/MetadataImageTest.java @@ -19,7 +19,6 @@ import org.apache.kafka.metadata.RecordTestUtils; import org.apache.kafka.raft.OffsetAndEpoch; -import org.apache.kafka.server.common.MetadataVersion; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; @@ -87,7 +86,7 @@ public void testImage2RoundTrip() throws Throwable { private void testToImageAndBack(MetadataImage image) throws Throwable { MockSnapshotConsumer writer = new MockSnapshotConsumer(); - image.write(writer, MetadataVersion.latest()); + image.write(writer); MetadataDelta delta = new MetadataDelta(MetadataImage.EMPTY); RecordTestUtils.replayAllBatches( delta, image.highestOffsetAndEpoch().offset, image.highestOffsetAndEpoch().epoch, writer.batches()); From db8224767457c79e8af31af0b4ddbd9988f41086 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Mon, 6 Jun 2022 08:15:45 +0200 Subject: [PATCH 11/18] address minor comments --- .../controller/ClusterControlManager.java | 63 ++++++++++--------- .../controller/ReplicationControlManager.java | 3 +- .../org/apache/kafka/image/ClusterDelta.java | 51 ++++++--------- .../kafka/metadata/BrokerRegistration.java | 21 ++++--- ...egistrationInControlledShutdownChange.java | 3 + 5 files changed, 71 insertions(+), 70 deletions(-) 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 b86ad8440a081..4cc31d3ec5ed6 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -432,40 +432,49 @@ public void replay(UnregisterBrokerRecord record) { } public void replay(FenceBrokerRecord record) { - replayRegistrationChange(record, record.id(), record.epoch(), - BrokerRegistrationFencingChange.UNFENCE, - BrokerRegistrationInControlledShutdownChange.NONE); + replayRegistrationChange( + record, + record.id(), + record.epoch(), + BrokerRegistrationFencingChange.UNFENCE.asBoolean(), + BrokerRegistrationInControlledShutdownChange.NONE.asBoolean() + ); } public void replay(UnfenceBrokerRecord record) { - replayRegistrationChange(record, record.id(), record.epoch(), - BrokerRegistrationFencingChange.FENCE, - BrokerRegistrationInControlledShutdownChange.NONE); + replayRegistrationChange( + record, + record.id(), + record.epoch(), + BrokerRegistrationFencingChange.FENCE.asBoolean(), + BrokerRegistrationInControlledShutdownChange.NONE.asBoolean() + ); } public void replay(BrokerRegistrationChangeRecord record) { - Optional fencingChange = - BrokerRegistrationFencingChange.fromValue(record.fenced()); - if (!fencingChange.isPresent()) { - throw new RuntimeException(String.format("Unable to replay %s: unknown " + - "value for fenced field: %d", record.toString(), record.fenced())); - } - Optional inControlledShutdownChange = - BrokerRegistrationInControlledShutdownChange.fromValue(record.inControlledShutdown()); - if (!inControlledShutdownChange.isPresent()) { - throw new RuntimeException(String.format("Unable to replay %s: unknown " + - "value for inControlledShutdown field: %d", record.toString(), record.inControlledShutdown())); - } - replayRegistrationChange(record, record.brokerId(), record.brokerEpoch(), - fencingChange.get(), inControlledShutdownChange.get()); + BrokerRegistrationFencingChange fencingChange = + BrokerRegistrationFencingChange.fromValue(record.fenced()).orElseThrow( + () -> new IllegalStateException(String.format("Unable to replay %s: unknown " + + "value for fenced field: %d", record, record.fenced()))); + BrokerRegistrationInControlledShutdownChange inControlledShutdownChange = + BrokerRegistrationInControlledShutdownChange.fromValue(record.inControlledShutdown()).orElseThrow( + () -> new IllegalStateException(String.format("Unable to replay %s: unknown " + + "value for inControlledShutdown field: %d", record, record.inControlledShutdown()))); + replayRegistrationChange( + record, + record.brokerId(), + record.brokerEpoch(), + fencingChange.asBoolean(), + inControlledShutdownChange.asBoolean() + ); } private void replayRegistrationChange( ApiMessage record, int brokerId, long brokerEpoch, - BrokerRegistrationFencingChange fencingChange, - BrokerRegistrationInControlledShutdownChange inControlledShutdownChange + Optional fencingChange, + Optional inControlledShutdownChange ) { BrokerRegistration curRegistration = brokerRegistrations.get(brokerId); if (curRegistration == null) { @@ -475,12 +484,10 @@ private void replayRegistrationChange( throw new RuntimeException(String.format("Unable to replay %s: no broker " + "registration with that epoch found", record.toString())); } else { - BrokerRegistration nextRegistration = curRegistration; - if (fencingChange != BrokerRegistrationFencingChange.NONE - || inControlledShutdownChange != BrokerRegistrationInControlledShutdownChange.NONE) { - nextRegistration = nextRegistration.cloneWith( - fencingChange.asBoolean(), inControlledShutdownChange.asBoolean()); - } + BrokerRegistration nextRegistration = curRegistration.maybeCloneWith( + fencingChange, + inControlledShutdownChange + ).orElse(curRegistration); if (!curRegistration.equals(nextRegistration)) { brokerRegistrations.put(brokerId, nextRegistration); updateMetrics(curRegistration, nextRegistration); diff --git a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java index 52bc28e56d019..b30ef79471834 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java @@ -1183,7 +1183,8 @@ void handleBrokerUnfenced(int brokerId, long brokerEpoch, List records) { - if (featureControl.metadataVersion().isInControlledShutdownStateSupported()) { + if (featureControl.metadataVersion().isInControlledShutdownStateSupported() + && !clusterControl.inControlledShutdown(brokerId)) { records.add(new ApiMessageAndVersion(new BrokerRegistrationChangeRecord(). setBrokerId(brokerId).setBrokerEpoch(brokerEpoch). setInControlledShutdown(BrokerRegistrationInControlledShutdownChange.IN_CONTROLLED_SHUTDOWN.value()), 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 14bb8620ac4fe..251a83f873ec3 100644 --- a/metadata/src/main/java/org/apache/kafka/image/ClusterDelta.java +++ b/metadata/src/main/java/org/apache/kafka/image/ClusterDelta.java @@ -93,46 +93,35 @@ private BrokerRegistration getBrokerOrThrow(int brokerId, long epoch, String act public void replay(FenceBrokerRecord record) { BrokerRegistration broker = getBrokerOrThrow(record.id(), record.epoch(), "fence"); - changedBrokers.put(record.id(), Optional.of(broker.cloneWithFencing(true))); + changedBrokers.put(record.id(), broker.maybeCloneWith( + BrokerRegistrationFencingChange.UNFENCE.asBoolean(), + Optional.empty() + )); } public void replay(UnfenceBrokerRecord record) { BrokerRegistration broker = getBrokerOrThrow(record.id(), record.epoch(), "unfence"); - changedBrokers.put(record.id(), Optional.of(broker.cloneWithFencing(false))); + changedBrokers.put(record.id(), broker.maybeCloneWith( + BrokerRegistrationFencingChange.FENCE.asBoolean(), + Optional.empty() + )); } public void replay(BrokerRegistrationChangeRecord record) { BrokerRegistration broker = getBrokerOrThrow(record.brokerId(), record.brokerEpoch(), "change"); - Optional fencingChange = - BrokerRegistrationFencingChange.fromValue(record.fenced()); - if (!fencingChange.isPresent()) { - throw new IllegalStateException(String.format("Unable to replay %s: unknown " + - "value for fenced field: %d", record.toString(), record.fenced())); - } - Optional inControlledShutdownChange = - BrokerRegistrationInControlledShutdownChange.fromValue(record.inControlledShutdown()); - if (!inControlledShutdownChange.isPresent()) { - throw new IllegalStateException(String.format("Unable to replay %s: unknown " + - "value for inControlledShutdown field: %d", record.toString(), record.inControlledShutdown())); - } - - replayRegistration(record.brokerId(), broker, fencingChange.get(), inControlledShutdownChange.get()); - } - - private void replayRegistration( - int brokerId, - BrokerRegistration broker, - BrokerRegistrationFencingChange fencingChange, - BrokerRegistrationInControlledShutdownChange inControlledShutdownChange - ) { - if (fencingChange != BrokerRegistrationFencingChange.NONE - || inControlledShutdownChange != BrokerRegistrationInControlledShutdownChange.NONE) { - changedBrokers.put(brokerId, Optional.of(broker.cloneWith( - fencingChange.asBoolean(), - inControlledShutdownChange.asBoolean() - ))); - } + BrokerRegistrationFencingChange fencingChange = + BrokerRegistrationFencingChange.fromValue(record.fenced()).orElseThrow( + () -> new IllegalStateException(String.format("Unable to replay %s: unknown " + + "value for fenced field: %d", record, record.fenced()))); + BrokerRegistrationInControlledShutdownChange inControlledShutdownChange = + BrokerRegistrationInControlledShutdownChange.fromValue(record.inControlledShutdown()).orElseThrow( + () -> new IllegalStateException(String.format("Unable to replay %s: unknown " + + "value for inControlledShutdown field: %d", record, record.inControlledShutdown()))); + changedBrokers.put(record.brokerId(), broker.maybeCloneWith( + fencingChange.asBoolean(), + inControlledShutdownChange.asBoolean() + )); } public ClusterImage apply() { diff --git a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java index 2c0a9aab67174..b79b100d4b0fa 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java @@ -235,24 +235,25 @@ public String toString() { return bld.toString(); } - public BrokerRegistration cloneWithFencing(boolean fencing) { - return new BrokerRegistration(id, epoch, incarnationId, listeners, - supportedFeatures, rack, fencing, inControlledShutdown); - } - - public BrokerRegistration cloneWith( + public Optional maybeCloneWith( Optional fencingChange, Optional inControlledShutdownChange ) { - return new BrokerRegistration( + boolean newFenced = fencingChange.orElse(fenced); + boolean newInControlledShutdownChange = inControlledShutdownChange.orElse(inControlledShutdown); + + if (newFenced == fenced && newInControlledShutdownChange == inControlledShutdown) + return Optional.empty(); + + return Optional.of(new BrokerRegistration( id, epoch, incarnationId, listeners, supportedFeatures, rack, - fencingChange.orElse(fenced), - inControlledShutdownChange.orElse(inControlledShutdown) - ); + newFenced, + newInControlledShutdownChange + )); } } diff --git a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistrationInControlledShutdownChange.java b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistrationInControlledShutdownChange.java index 0cd4c3959a249..39f8abf595e7a 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistrationInControlledShutdownChange.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistrationInControlledShutdownChange.java @@ -24,6 +24,9 @@ import java.util.stream.Collectors; public enum BrokerRegistrationInControlledShutdownChange { + // Note that Optional.of(true) is not a valid state change here. The only + // way to leave the in controlled shutdown state is by registering the + // broker with a new incarnation id. NONE(0, Optional.empty()), IN_CONTROLLED_SHUTDOWN(1, Optional.of(true)); From 13e7c0ac4232407f2a129f34e75266f97be1fc84 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Mon, 6 Jun 2022 08:19:57 +0200 Subject: [PATCH 12/18] add javadoc --- .../kafka/controller/ClusterControlManager.java | 12 ++++++++++++ 1 file changed, 12 insertions(+) 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 4cc31d3ec5ed6..910e1ce779ebd 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -541,18 +541,30 @@ Iterator usableBrokers() { id -> brokerRegistrations.get(id).rack()); } + /** + * Returns true if the broker is in fenced state; Returns false if it is + * not or if it does not exist. + */ public boolean unfenced(int brokerId) { BrokerRegistration registration = brokerRegistrations.get(brokerId); if (registration == null) return false; return !registration.fenced(); } + /** + * Returns true if the broker is in controlled shutdown state; Returns false + * if it is not or if it does not exist. + */ public boolean inControlledShutdown(int brokerId) { BrokerRegistration registration = brokerRegistrations.get(brokerId); if (registration == null) return false; return registration.inControlledShutdown(); } + /** + * Returns true if the broker is active. Active means not fenced nor in controlled + * shutdown; Returns false if it is not active or if it does not exist. + */ public boolean active(int brokerId) { BrokerRegistration registration = brokerRegistrations.get(brokerId); if (registration == null) return false; From bb2bb7c2bfab5b50a94432e72dc5fef441e32680 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Mon, 6 Jun 2022 10:06:41 +0200 Subject: [PATCH 13/18] address comments --- .../controller/ReplicationControlManager.java | 20 ++++- .../ReplicationControlManagerTest.java | 79 +++++++++++++++---- 2 files changed, 80 insertions(+), 19 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java index b30ef79471834..8794a183649aa 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java @@ -692,8 +692,14 @@ private ApiError createTopic(CreatableTopic topic, List replicas = partitions.get(partitionId); List isr = replicas.stream(). filter(clusterControl::active).collect(Collectors.toList()); - // We need to have at least one replica in the ISR. - if (isr.isEmpty()) isr.add(replicas.get(0)); + // If the ISR is empty, it means that all brokers are fenced or + // in controlled shutdown. To be consistent with the replica placer, + // we reject the create topic request with INVALID_REPLICATION_FACTOR. + if (isr.isEmpty()) { + return new ApiError(Errors.INVALID_REPLICATION_FACTOR, + "Unable to replicate the partition " + replicationFactor + + " time(s): All brokers are currently fenced or in controlled shutdown."); + } newParts.put(partitionId, new PartitionRegistration( Replicas.toArray(replicas), @@ -1503,8 +1509,14 @@ void createPartitions(CreatePartitionsTopic topic, List replicas = placements.get(i); List isr = isrs.get(i).stream(). filter(clusterControl::active).collect(Collectors.toList()); - // We need to have at least one replica in the ISR. - if (isr.isEmpty()) isr.add(replicas.get(0)); + // If the ISR is empty, it means that all brokers are fenced or + // in controlled shutdown. To be consistent with the replica placer, + // we reject the create topic request with INVALID_REPLICATION_FACTOR. + if (isr.isEmpty()) { + throw new InvalidReplicationFactorException( + "Unable to replicate the partition " + replicationFactor + + " time(s): All brokers are currently fenced or in controlled shutdown."); + } records.add(new ApiMessageAndVersion(new PartitionRecord(). setPartitionId(partitionId). setTopicId(topicId). diff --git a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java index 1dee7bfad5e70..fc63f977dd3f5 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java @@ -114,6 +114,7 @@ import static org.apache.kafka.common.protocol.Errors.ELIGIBLE_LEADERS_NOT_AVAILABLE; import static org.apache.kafka.common.protocol.Errors.FENCED_LEADER_EPOCH; import static org.apache.kafka.common.protocol.Errors.INVALID_PARTITIONS; +import static org.apache.kafka.common.protocol.Errors.INVALID_REPLICATION_FACTOR; import static org.apache.kafka.common.protocol.Errors.INVALID_REPLICA_ASSIGNMENT; import static org.apache.kafka.common.protocol.Errors.INVALID_TOPIC_EXCEPTION; import static org.apache.kafka.common.protocol.Errors.NONE; @@ -438,38 +439,53 @@ public void testCreateTopics() throws Exception { CreateTopicsRequestData request = new CreateTopicsRequestData(); request.topics().add(new CreatableTopic().setName("foo"). setNumPartitions(-1).setReplicationFactor((short) -1)); + ControllerResult result = replicationControl.createTopics(request, Collections.singleton("foo")); CreateTopicsResponseData expectedResponse = new CreateTopicsResponseData(); expectedResponse.topics().add(new CreatableTopicResult().setName("foo"). - setErrorCode(Errors.INVALID_REPLICATION_FACTOR.code()). + setErrorCode(INVALID_REPLICATION_FACTOR.code()). setErrorMessage("Unable to replicate the partition 3 time(s): All " + "brokers are currently fenced.")); assertEquals(expectedResponse, result.response()); ctx.registerBrokers(0, 1, 2); - ctx.unfenceBrokers(0, 1, 2); + ctx.unfenceBrokers(0); + ctx.inControlledShutdownBrokers(0); + ControllerResult result2 = replicationControl.createTopics(request, Collections.singleton("foo")); CreateTopicsResponseData expectedResponse2 = new CreateTopicsResponseData(); expectedResponse2.topics().add(new CreatableTopicResult().setName("foo"). + setErrorCode(INVALID_REPLICATION_FACTOR.code()). + setErrorMessage("Unable to replicate the partition 3 time(s): All " + + "brokers are currently fenced or in controlled shutdown.")); + assertEquals(expectedResponse2, result2.response()); + + ctx.registerBrokers(0, 1, 2); + ctx.unfenceBrokers(0, 1, 2); + + ControllerResult result3 = + replicationControl.createTopics(request, Collections.singleton("foo")); + CreateTopicsResponseData expectedResponse3 = new CreateTopicsResponseData(); + expectedResponse3.topics().add(new CreatableTopicResult().setName("foo"). setNumPartitions(1).setReplicationFactor((short) 3). setErrorMessage(null).setErrorCode((short) 0). - setTopicId(result2.response().topics().find("foo").topicId())); - assertEquals(expectedResponse2, result2.response()); - ctx.replay(result2.records()); + setTopicId(result3.response().topics().find("foo").topicId())); + assertEquals(expectedResponse3, result3.response()); + ctx.replay(result3.records()); assertEquals(new PartitionRegistration(new int[] {1, 2, 0}, new int[] {1, 2, 0}, Replicas.NONE, Replicas.NONE, 1, LeaderRecoveryState.RECOVERED, 0, 0), replicationControl.getPartition( - ((TopicRecord) result2.records().get(0).message()).topicId(), 0)); - ControllerResult result3 = + ((TopicRecord) result3.records().get(0).message()).topicId(), 0)); + ControllerResult result4 = replicationControl.createTopics(request, Collections.singleton("foo")); - CreateTopicsResponseData expectedResponse3 = new CreateTopicsResponseData(); - expectedResponse3.topics().add(new CreatableTopicResult().setName("foo"). + CreateTopicsResponseData expectedResponse4 = new CreateTopicsResponseData(); + expectedResponse4.topics().add(new CreatableTopicResult().setName("foo"). setErrorCode(Errors.TOPIC_ALREADY_EXISTS.code()). setErrorMessage("Topic 'foo' already exists.")); - assertEquals(expectedResponse3, result3.response()); - Uuid fooId = result2.response().topics().find("foo").topicId(); + assertEquals(expectedResponse4, result4.response()); + Uuid fooId = result3.response().topics().find("foo").topicId(); RecordTestUtils.assertBatchIteratorContains(asList( asList(new ApiMessageAndVersion(new PartitionRecord(). setPartitionId(0).setTopicId(fooId). @@ -482,7 +498,7 @@ public void testCreateTopics() throws Exception { } @Test - public void testCreateTopicsInvariants() throws Exception { + public void testCreateTopicsISRInvariants() throws Exception { ReplicationControlTestContext ctx = new ReplicationControlTestContext(); ReplicationControlManager replicationControl = ctx.replicationControl; @@ -634,7 +650,7 @@ public void testInvalidCreateTopicsWithValidateOnlyFlag() throws Exception { assertEquals(0, result.records().size()); CreateTopicsResponseData expectedResponse = new CreateTopicsResponseData(); expectedResponse.topics().add(new CreatableTopicResult().setName("foo"). - setErrorCode(Errors.INVALID_REPLICATION_FACTOR.code()). + setErrorCode(INVALID_REPLICATION_FACTOR.code()). setErrorMessage("Unable to replicate the partition 4 time(s): The target " + "replication factor of 4 cannot be reached because only 3 broker(s) " + "are registered.")); @@ -1091,7 +1107,6 @@ public void testDeleteTopics() throws Exception { assertEmptyTopicConfigs(ctx, "foo"); } - @Test public void testCreatePartitions() throws Exception { ReplicationControlTestContext ctx = new ReplicationControlTestContext(); @@ -1178,7 +1193,41 @@ public void testCreatePartitions() throws Exception { } @Test - public void testCreatePartitionsInvariants() throws Exception { + public void testCreatePartitionsFailsWhenAllBrokersAreFencedOrInControlledShutdown() throws Exception { + ReplicationControlTestContext ctx = new ReplicationControlTestContext(); + ReplicationControlManager replicationControl = ctx.replicationControl; + CreateTopicsRequestData request = new CreateTopicsRequestData(); + request.topics().add(new CreatableTopic().setName("foo"). + setNumPartitions(1).setReplicationFactor((short) 2)); + + ctx.registerBrokers(0, 1); + ctx.unfenceBrokers(0, 1); + + ControllerResult createTopicResult = replicationControl. + createTopics(request, new HashSet<>(Arrays.asList("foo"))); + ctx.replay(createTopicResult.records()); + + ctx.registerBrokers(0, 1); + ctx.unfenceBrokers(0); + ctx.inControlledShutdownBrokers(0); + + List topics = new ArrayList<>(); + topics.add(new CreatePartitionsTopic(). + setName("foo").setCount(2).setAssignments(null)); + ControllerResult> createPartitionsResult = + replicationControl.createPartitions(topics); + + assertEquals( + asList(new CreatePartitionsTopicResult(). + setName("foo"). + setErrorCode(INVALID_REPLICATION_FACTOR.code()). + setErrorMessage("Unable to replicate the partition 2 time(s): All " + + "brokers are currently fenced or in controlled shutdown.")), + createPartitionsResult.response()); + } + + @Test + public void testCreatePartitionsISRInvariants() throws Exception { ReplicationControlTestContext ctx = new ReplicationControlTestContext(); ReplicationControlManager replicationControl = ctx.replicationControl; From 2761eef1a2c44c5e4304872838eb00cb30cc344f Mon Sep 17 00:00:00 2001 From: David Jacot Date: Mon, 6 Jun 2022 12:57:16 +0200 Subject: [PATCH 14/18] use MetadataVersion.IBP_3_0_IV1 by default --- .../src/main/java/org/apache/kafka/image/MetadataImage.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/image/MetadataImage.java b/metadata/src/main/java/org/apache/kafka/image/MetadataImage.java index d1aa3c6e8a938..e3cd94a0cb5bf 100644 --- a/metadata/src/main/java/org/apache/kafka/image/MetadataImage.java +++ b/metadata/src/main/java/org/apache/kafka/image/MetadataImage.java @@ -121,11 +121,11 @@ public AclsImage acls() { } public void write(Consumer> out) { - // We use the latest metadata version if this image does not have - // a specific version set. + // We use the minimum KRaft metadata version if this image does + // not have a specific version set. MetadataVersion metadataVersion = features.metadataVersion(); if (metadataVersion.equals(MetadataVersion.UNINITIALIZED)) { - metadataVersion = MetadataVersion.latest(); + metadataVersion = MetadataVersion.IBP_3_0_IV1; } // Features should be written out first so we can include the metadata.version at the beginning of the // snapshot From 87446cd737dad59cda90994baa03419adb7a812d Mon Sep 17 00:00:00 2001 From: David Jacot Date: Tue, 7 Jun 2022 10:14:08 +0200 Subject: [PATCH 15/18] refactor and address comments --- .../controller/ClusterControlManager.java | 25 ++++++++++--------- .../org/apache/kafka/image/ClusterDelta.java | 21 +++++++++------- .../kafka/metadata/BrokerRegistration.java | 13 +++++----- .../controller/ClusterControlManagerTest.java | 13 +++++++--- .../ReplicationControlManagerTest.java | 3 +-- .../kafka/server/common/MetadataVersion.java | 8 ++++++ .../server/common/MetadataVersionTest.java | 19 +++++++++++++- 7 files changed, 68 insertions(+), 34 deletions(-) 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 910e1ce779ebd..cbb161eeecea2 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -309,13 +309,6 @@ Set fencedBrokerIds() { .collect(Collectors.toSet()); } - private short registerBrokerRecordVersion() { - if (featureControl.metadataVersion().isInControlledShutdownStateSupported()) - return (short) 1; - else - return (short) 0; - } - /** * Process an incoming broker registration request. */ @@ -376,7 +369,8 @@ public ControllerResult registerBroker( heartbeatManager.register(brokerId, record.fenced()); List records = new ArrayList<>(); - records.add(new ApiMessageAndVersion(record, registerBrokerRecordVersion())); + records.add(new ApiMessageAndVersion(record, featureControl.metadataVersion(). + registerBrokerRecordVersion())); return ControllerResult.atomicOf(records, new BrokerRegistrationReply(brokerEpoch)); } @@ -484,10 +478,10 @@ private void replayRegistrationChange( throw new RuntimeException(String.format("Unable to replay %s: no broker " + "registration with that epoch found", record.toString())); } else { - BrokerRegistration nextRegistration = curRegistration.maybeCloneWith( + BrokerRegistration nextRegistration = curRegistration.cloneWith( fencingChange, inControlledShutdownChange - ).orElse(curRegistration); + ); if (!curRegistration.equals(nextRegistration)) { brokerRegistrations.put(brokerId, nextRegistration); updateMetrics(curRegistration, nextRegistration); @@ -600,9 +594,15 @@ public void addReadyBrokersFuture(CompletableFuture future, int minBrokers class ClusterControlIterator implements Iterator> { private final Iterator> iterator; + private final MetadataVersion metadataVersion; ClusterControlIterator(long epoch) { this.iterator = brokerRegistrations.entrySet(epoch).iterator(); + if (featureControl.metadataVersion().equals(MetadataVersion.UNINITIALIZED)) { + this.metadataVersion = MetadataVersion.IBP_3_0_IV1; + } else { + this.metadataVersion = featureControl.metadataVersion(); + } } @Override @@ -637,10 +637,11 @@ public List next() { setFeatures(features). setRack(registration.rack().orElse(null)). setFenced(registration.fenced()); - if (featureControl.metadataVersion().isInControlledShutdownStateSupported()) { + if (metadataVersion.isInControlledShutdownStateSupported()) { record.setInControlledShutdown(registration.inControlledShutdown()); } - return singletonList(new ApiMessageAndVersion(record, registerBrokerRecordVersion())); + return singletonList(new ApiMessageAndVersion(record, + metadataVersion.registerBrokerRecordVersion())); } } 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 251a83f873ec3..110e68fb89a29 100644 --- a/metadata/src/main/java/org/apache/kafka/image/ClusterDelta.java +++ b/metadata/src/main/java/org/apache/kafka/image/ClusterDelta.java @@ -92,23 +92,23 @@ private BrokerRegistration getBrokerOrThrow(int brokerId, long epoch, String act } public void replay(FenceBrokerRecord record) { - BrokerRegistration broker = getBrokerOrThrow(record.id(), record.epoch(), "fence"); - changedBrokers.put(record.id(), broker.maybeCloneWith( + BrokerRegistration curRegistration = getBrokerOrThrow(record.id(), record.epoch(), "fence"); + changedBrokers.put(record.id(), Optional.of(curRegistration.cloneWith( BrokerRegistrationFencingChange.UNFENCE.asBoolean(), Optional.empty() - )); + ))); } public void replay(UnfenceBrokerRecord record) { - BrokerRegistration broker = getBrokerOrThrow(record.id(), record.epoch(), "unfence"); - changedBrokers.put(record.id(), broker.maybeCloneWith( + BrokerRegistration curRegistration = getBrokerOrThrow(record.id(), record.epoch(), "unfence"); + changedBrokers.put(record.id(), Optional.of(curRegistration.cloneWith( BrokerRegistrationFencingChange.FENCE.asBoolean(), Optional.empty() - )); + ))); } public void replay(BrokerRegistrationChangeRecord record) { - BrokerRegistration broker = + BrokerRegistration curRegistration = getBrokerOrThrow(record.brokerId(), record.brokerEpoch(), "change"); BrokerRegistrationFencingChange fencingChange = BrokerRegistrationFencingChange.fromValue(record.fenced()).orElseThrow( @@ -118,10 +118,13 @@ public void replay(BrokerRegistrationChangeRecord record) { BrokerRegistrationInControlledShutdownChange.fromValue(record.inControlledShutdown()).orElseThrow( () -> new IllegalStateException(String.format("Unable to replay %s: unknown " + "value for inControlledShutdown field: %d", record, record.inControlledShutdown()))); - changedBrokers.put(record.brokerId(), broker.maybeCloneWith( + BrokerRegistration nextRegistration = curRegistration.cloneWith( fencingChange.asBoolean(), inControlledShutdownChange.asBoolean() - )); + ); + if (!curRegistration.equals(nextRegistration)) { + changedBrokers.put(record.brokerId(), Optional.of(nextRegistration)); + } } public ClusterImage apply() { diff --git a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java index b79b100d4b0fa..d1d345506530b 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/BrokerRegistration.java @@ -188,9 +188,8 @@ public ApiMessageAndVersion toRecord(MetadataVersion metadataVersion) { setMaxSupportedVersion(entry.getValue().max())); } - short version = metadataVersion.isInControlledShutdownStateSupported() ? - (short) 1 : (short) 0; - return new ApiMessageAndVersion(registrationRecord, version); + return new ApiMessageAndVersion(registrationRecord, + metadataVersion.registerBrokerRecordVersion()); } @Override @@ -235,7 +234,7 @@ public String toString() { return bld.toString(); } - public Optional maybeCloneWith( + public BrokerRegistration cloneWith( Optional fencingChange, Optional inControlledShutdownChange ) { @@ -243,9 +242,9 @@ public Optional maybeCloneWith( boolean newInControlledShutdownChange = inControlledShutdownChange.orElse(inControlledShutdown); if (newFenced == fenced && newInControlledShutdownChange == inControlledShutdown) - return Optional.empty(); + return this; - return Optional.of(new BrokerRegistration( + return new BrokerRegistration( id, epoch, incarnationId, @@ -254,6 +253,6 @@ public Optional maybeCloneWith( rack, newFenced, newInControlledShutdownChange - )); + ); } } 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 d5da9d921b5b4..c476083efb7ab 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ClusterControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ClusterControlManagerTest.java @@ -59,7 +59,6 @@ import org.junit.jupiter.params.provider.ValueSource; import static org.apache.kafka.server.common.MetadataVersion.IBP_3_3_IV2; -import static org.apache.kafka.server.common.MetadataVersion.IBP_3_3_IV3; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; @@ -269,7 +268,7 @@ public void testRegisterBrokerRecordVersion(MetadataVersion metadataVersion) { 123L, new FinalizedControllerFeatures(Collections.emptyMap(), 456L)); - short expectedVersion = metadataVersion.isAtLeast(IBP_3_3_IV3) ? (short) 1 : (short) 0; + short expectedVersion = metadataVersion.registerBrokerRecordVersion(); assertEquals( Arrays.asList(new ApiMessageAndVersion(new RegisterBrokerRecord(). @@ -402,7 +401,14 @@ public void testIterator(MetadataVersion metadataVersion) throws Exception { new UnfenceBrokerRecord().setId(i).setEpoch(100); clusterControl.replay(unfenceBrokerRecord); } - short expectedVersion = metadataVersion.isAtLeast(IBP_3_3_IV3) ? (short) 1 : (short) 0; + BrokerRegistrationChangeRecord registrationChangeRecord = + new BrokerRegistrationChangeRecord(). + setBrokerId(0). + setBrokerEpoch(100). + setInControlledShutdown(BrokerRegistrationInControlledShutdownChange. + IN_CONTROLLED_SHUTDOWN.value()); + clusterControl.replay(registrationChangeRecord); + short expectedVersion = metadataVersion.registerBrokerRecordVersion(); RecordTestUtils.assertBatchIteratorContains(Arrays.asList( Arrays.asList(new ApiMessageAndVersion(new RegisterBrokerRecord(). setBrokerEpoch(100).setBrokerId(0).setRack(null). @@ -411,6 +417,7 @@ public void testIterator(MetadataVersion metadataVersion) throws Exception { setPort((short) 9092). setName("PLAINTEXT"). setHost("example.com")).iterator())). + setInControlledShutdown(metadataVersion.isInControlledShutdownStateSupported()). setFenced(false), expectedVersion)), Arrays.asList(new ApiMessageAndVersion(new RegisterBrokerRecord(). setBrokerEpoch(100).setBrokerId(1).setRack(null). diff --git a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java index fc63f977dd3f5..764ba7fdd5cd9 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java @@ -124,7 +124,6 @@ import static org.apache.kafka.common.protocol.Errors.UNKNOWN_TOPIC_ID; import static org.apache.kafka.common.protocol.Errors.UNKNOWN_TOPIC_OR_PARTITION; import static org.apache.kafka.metadata.LeaderConstants.NO_LEADER; -import static org.apache.kafka.server.common.MetadataVersion.IBP_3_3_IV3; import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; @@ -2017,7 +2016,7 @@ public void testProcessBrokerHeartbeatInControlledShutdown(MetadataVersion metad List expectedRecords = new ArrayList<>(); - if (metadataVersion.isAtLeast(IBP_3_3_IV3)) { + if (metadataVersion.isInControlledShutdownStateSupported()) { expectedRecords.add(new ApiMessageAndVersion( new BrokerRegistrationChangeRecord() .setBrokerEpoch(100) 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 cdf69f993cbd5..aa1a416c2943a 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 @@ -250,6 +250,14 @@ public boolean isInControlledShutdownStateSupported() { return this.isAtLeast(IBP_3_3_IV3); } + public short registerBrokerRecordVersion() { + if (isInControlledShutdownStateSupported()) { + return (short) 1; + } else { + return (short) 0; + } + } + private static final Map IBP_VERSIONS; static { { diff --git a/server-common/src/test/java/org/apache/kafka/server/common/MetadataVersionTest.java b/server-common/src/test/java/org/apache/kafka/server/common/MetadataVersionTest.java index ec038383caad4..2f96d5fd04fb8 100644 --- a/server-common/src/test/java/org/apache/kafka/server/common/MetadataVersionTest.java +++ b/server-common/src/test/java/org/apache/kafka/server/common/MetadataVersionTest.java @@ -18,9 +18,11 @@ package org.apache.kafka.server.common; import org.apache.kafka.common.record.RecordVersion; -import org.junit.jupiter.api.Test; import java.util.Arrays; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; import static org.apache.kafka.server.common.MetadataVersion.IBP_0_10_0_IV0; import static org.apache.kafka.server.common.MetadataVersion.IBP_0_10_0_IV1; @@ -322,4 +324,19 @@ public void testKRaftVersions() { } } } + + @ParameterizedTest + @EnumSource(value = MetadataVersion.class) + public void testIsInControlledShutdownStateSupported(MetadataVersion metadataVersion) { + assertEquals(metadataVersion.isAtLeast(IBP_3_3_IV3), + metadataVersion.isInControlledShutdownStateSupported()); + } + + @ParameterizedTest + @EnumSource(value = MetadataVersion.class) + public void testRegisterBrokerRecordVersion(MetadataVersion metadataVersion) { + short expectedVersion = metadataVersion.isAtLeast(IBP_3_3_IV3) ? + (short) 1 : (short) 0; + assertEquals(expectedVersion, metadataVersion.registerBrokerRecordVersion()); + } } From 1450f431b2cabfe2d540010259add86f1fe9defe Mon Sep 17 00:00:00 2001 From: David Jacot Date: Tue, 7 Jun 2022 10:23:00 +0200 Subject: [PATCH 16/18] don't automagically instanciate FeatureControlManager in ClusterControlManager.Builder --- .../controller/ClusterControlManager.java | 10 +--- .../controller/ClusterControlManagerTest.java | 48 +++++++++++++++++++ .../ProducerIdControlManagerTest.java | 12 +++++ .../ReplicationControlManagerTest.java | 8 ++++ 4 files changed, 69 insertions(+), 9 deletions(-) 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 cbb161eeecea2..21bd01c3b92f7 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -17,7 +17,6 @@ package org.apache.kafka.controller; -import org.apache.kafka.clients.ApiVersions; import org.apache.kafka.common.Endpoint; import org.apache.kafka.common.Uuid; import org.apache.kafka.common.errors.DuplicateBrokerRegistrationException; @@ -146,14 +145,7 @@ ClusterControlManager build() { throw new RuntimeException("You must specify ControllerMetrics"); } if (featureControl == null) { - featureControl = new FeatureControlManager.Builder(). - setLogContext(logContext). - setSnapshotRegistry(snapshotRegistry). - setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), - QuorumFeatures.defaultFeatureMap(), - singletonList(0))). - setMetadataVersion(MetadataVersion.latest()). - build(); + throw new RuntimeException("You must specify FeatureControlManager"); } return new ClusterControlManager(logContext, clusterId, 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 c476083efb7ab..cd44e2e678cdb 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ClusterControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ClusterControlManagerTest.java @@ -73,11 +73,19 @@ public void testReplay(MetadataVersion metadataVersion) { MockTime time = new MockTime(0, 0, 0); SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); + FeatureControlManager featureControl = new FeatureControlManager.Builder(). + setSnapshotRegistry(snapshotRegistry). + setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), + QuorumFeatures.defaultFeatureMap(), + Collections.singletonList(0))). + setMetadataVersion(MetadataVersion.latest()). + build(); ClusterControlManager clusterControl = new ClusterControlManager.Builder(). setTime(time). setSnapshotRegistry(snapshotRegistry). setSessionTimeoutNs(1000). setControllerMetrics(new MockControllerMetrics()). + setFeatureControlManager(featureControl). build(); clusterControl.activate(); assertFalse(clusterControl.unfenced(0)); @@ -127,12 +135,20 @@ public void testReplayRegisterBrokerRecord() { MockTime time = new MockTime(0, 0, 0); SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); + FeatureControlManager featureControl = new FeatureControlManager.Builder(). + setSnapshotRegistry(snapshotRegistry). + setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), + QuorumFeatures.defaultFeatureMap(), + Collections.singletonList(0))). + setMetadataVersion(MetadataVersion.latest()). + build(); ClusterControlManager clusterControl = new ClusterControlManager.Builder(). setClusterId("fPZv1VBsRFmnlRvmGcOW9w"). setTime(time). setSnapshotRegistry(snapshotRegistry). setSessionTimeoutNs(1000). setControllerMetrics(new MockControllerMetrics()). + setFeatureControlManager(featureControl). build(); assertFalse(clusterControl.unfenced(0)); @@ -172,12 +188,20 @@ public void testReplayBrokerRegistrationChangeRecord() { MockTime time = new MockTime(0, 0, 0); SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); + FeatureControlManager featureControl = new FeatureControlManager.Builder(). + setSnapshotRegistry(snapshotRegistry). + setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), + QuorumFeatures.defaultFeatureMap(), + Collections.singletonList(0))). + setMetadataVersion(MetadataVersion.latest()). + build(); ClusterControlManager clusterControl = new ClusterControlManager.Builder(). setClusterId("fPZv1VBsRFmnlRvmGcOW9w"). setTime(time). setSnapshotRegistry(snapshotRegistry). setSessionTimeoutNs(1000). setControllerMetrics(new MockControllerMetrics()). + setFeatureControlManager(featureControl). build(); assertFalse(clusterControl.unfenced(0)); @@ -220,12 +244,20 @@ public void testReplayBrokerRegistrationChangeRecord() { @Test public void testRegistrationWithIncorrectClusterId() throws Exception { SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); + FeatureControlManager featureControl = new FeatureControlManager.Builder(). + setSnapshotRegistry(snapshotRegistry). + setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), + QuorumFeatures.defaultFeatureMap(), + Collections.singletonList(0))). + setMetadataVersion(MetadataVersion.latest()). + build(); ClusterControlManager clusterControl = new ClusterControlManager.Builder(). setClusterId("fPZv1VBsRFmnlRvmGcOW9w"). setTime(new MockTime(0, 0, 0)). setSnapshotRegistry(snapshotRegistry). setSessionTimeoutNs(1000). setControllerMetrics(new MockControllerMetrics()). + setFeatureControlManager(featureControl). build(); clusterControl.activate(); assertThrows(InconsistentClusterIdException.class, () -> @@ -294,11 +326,19 @@ public void testUnregister() throws Exception { setName("PLAINTEXT"). setHost("example.com")); SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); + FeatureControlManager featureControl = new FeatureControlManager.Builder(). + setSnapshotRegistry(snapshotRegistry). + setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), + QuorumFeatures.defaultFeatureMap(), + Collections.singletonList(0))). + setMetadataVersion(MetadataVersion.latest()). + build(); ClusterControlManager clusterControl = new ClusterControlManager.Builder(). setTime(new MockTime(0, 0, 0)). setSnapshotRegistry(snapshotRegistry). setSessionTimeoutNs(1000). setControllerMetrics(new MockControllerMetrics()). + setFeatureControlManager(featureControl). build(); clusterControl.activate(); clusterControl.replay(brokerRecord); @@ -319,11 +359,19 @@ public void testUnregister() throws Exception { public void testPlaceReplicas(int numUsableBrokers) throws Exception { MockTime time = new MockTime(0, 0, 0); SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); + FeatureControlManager featureControl = new FeatureControlManager.Builder(). + setSnapshotRegistry(snapshotRegistry). + setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), + QuorumFeatures.defaultFeatureMap(), + Collections.singletonList(0))). + setMetadataVersion(MetadataVersion.latest()). + build(); ClusterControlManager clusterControl = new ClusterControlManager.Builder(). setTime(time). setSnapshotRegistry(snapshotRegistry). setSessionTimeoutNs(1000). setControllerMetrics(new MockControllerMetrics()). + setFeatureControlManager(featureControl). build(); clusterControl.activate(); for (int i = 0; i < numUsableBrokers; i++) { diff --git a/metadata/src/test/java/org/apache/kafka/controller/ProducerIdControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/ProducerIdControlManagerTest.java index 284b8f7c1673c..ccdd3a5b2331f 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ProducerIdControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ProducerIdControlManagerTest.java @@ -17,6 +17,8 @@ package org.apache.kafka.controller; +import java.util.Collections; +import org.apache.kafka.clients.ApiVersions; import org.apache.kafka.common.errors.StaleBrokerEpochException; import org.apache.kafka.common.errors.UnknownServerException; import org.apache.kafka.common.metadata.ProducerIdsRecord; @@ -25,6 +27,7 @@ import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.common.utils.MockTime; import org.apache.kafka.server.common.ApiMessageAndVersion; +import org.apache.kafka.server.common.MetadataVersion; import org.apache.kafka.server.common.ProducerIdsBlock; import org.apache.kafka.timeline.SnapshotRegistry; import org.junit.jupiter.api.BeforeEach; @@ -42,6 +45,7 @@ public class ProducerIdControlManagerTest { private SnapshotRegistry snapshotRegistry; + private FeatureControlManager featureControl; private ClusterControlManager clusterControl; private ProducerIdControlManager producerIdControlManager; @@ -49,11 +53,19 @@ public class ProducerIdControlManagerTest { public void setUp() { final MockTime time = new MockTime(); snapshotRegistry = new SnapshotRegistry(new LogContext()); + featureControl = new FeatureControlManager.Builder(). + setSnapshotRegistry(snapshotRegistry). + setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), + QuorumFeatures.defaultFeatureMap(), + Collections.singletonList(0))). + setMetadataVersion(MetadataVersion.latest()). + build(); clusterControl = new ClusterControlManager.Builder(). setTime(time). setSnapshotRegistry(snapshotRegistry). setSessionTimeoutNs(1000). setControllerMetrics(new MockControllerMetrics()). + setFeatureControlManager(featureControl). build(); clusterControl.activate(); diff --git a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java index 764ba7fdd5cd9..59b5488f6a104 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java @@ -145,6 +145,13 @@ private static class ReplicationControlTestContext { final MockTime time = new MockTime(); final MockRandom random = new MockRandom(); final ControllerMetrics metrics = new MockControllerMetrics(); + final FeatureControlManager featureControl = new FeatureControlManager.Builder(). + setSnapshotRegistry(snapshotRegistry). + setQuorumFeatures(new QuorumFeatures(0, new ApiVersions(), + QuorumFeatures.defaultFeatureMap(), + Collections.singletonList(0))). + setMetadataVersion(MetadataVersion.latest()). + build(); final ClusterControlManager clusterControl = new ClusterControlManager.Builder(). setLogContext(logContext). setTime(time). @@ -152,6 +159,7 @@ private static class ReplicationControlTestContext { setSessionTimeoutNs(TimeUnit.MILLISECONDS.convert(BROKER_SESSION_TIMEOUT_MS, TimeUnit.NANOSECONDS)). setReplicaPlacer(new StripedReplicaPlacer(random)). setControllerMetrics(metrics). + setFeatureControlManager(featureControl). build(); final ConfigurationControlManager configurationControl = new ConfigurationControlManager.Builder(). setSnapshotRegistry(snapshotRegistry). From 4bd208c3d371453d17422960a85a1ad15ca78f27 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Tue, 7 Jun 2022 10:50:51 +0200 Subject: [PATCH 17/18] refactor --- .../kafka/controller/QuorumController.java | 7 ++++++- .../controller/ReplicationControlManager.java | 5 ----- .../kafka/controller/QuorumControllerTest.java | 16 ++++++++-------- 3 files changed, 14 insertions(+), 14 deletions(-) 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 ea9b7205e5050..bf4632d3629fe 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -622,11 +622,16 @@ public String toString() { } } - // VisibleForTesting + // Visible for testing ReplicationControlManager replicationControl() { return replicationControl; } + // Visible for testing + ClusterControlManager clusterControl() { + return clusterControl; + } + CompletableFuture appendReadEvent( String name, OptionalLong deadlineNs, diff --git a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java index 8794a183649aa..1c81954617ad8 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java @@ -1410,11 +1410,6 @@ ControllerResult maybeBalancePartitionLeaders() { return ControllerResult.of(records, rescheduleImmidiately); } - // Visible for testing - Boolean isBrokerUnfenced(int brokerId) { - return clusterControl.unfenced(brokerId); - } - ControllerResult> createPartitions(List topics) { List records = new ArrayList<>(); diff --git a/metadata/src/test/java/org/apache/kafka/controller/QuorumControllerTest.java b/metadata/src/test/java/org/apache/kafka/controller/QuorumControllerTest.java index d429a13b0ff48..87afbf4199f7b 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/QuorumControllerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/QuorumControllerTest.java @@ -234,7 +234,7 @@ public void testFenceMultipleBrokers() throws Throwable { // Brokers are only registered and should still be fenced allBrokers.forEach(brokerId -> { - assertFalse(active.replicationControl().isBrokerUnfenced(brokerId), + assertFalse(active.clusterControl().unfenced(brokerId), "Broker " + brokerId + " should have been fenced"); }); @@ -254,7 +254,7 @@ public void testFenceMultipleBrokers() throws Throwable { TestUtils.waitForCondition(() -> { sendBrokerheartbeat(active, brokersToKeepUnfenced, brokerEpochs); for (Integer brokerId : brokersToFence) { - if (active.replicationControl().isBrokerUnfenced(brokerId)) { + if (active.clusterControl().unfenced(brokerId)) { return false; } } @@ -268,11 +268,11 @@ public void testFenceMultipleBrokers() throws Throwable { // At this point only the brokers we want fenced should be fenced. brokersToKeepUnfenced.forEach(brokerId -> { - assertTrue(active.replicationControl().isBrokerUnfenced(brokerId), + assertTrue(active.clusterControl().unfenced(brokerId), "Broker " + brokerId + " should have been unfenced"); }); brokersToFence.forEach(brokerId -> { - assertFalse(active.replicationControl().isBrokerUnfenced(brokerId), + assertFalse(active.clusterControl().unfenced(brokerId), "Broker " + brokerId + " should have been fenced"); }); @@ -326,7 +326,7 @@ public void testBalancePartitionLeaders() throws Throwable { // Brokers are only registered and should still be fenced allBrokers.forEach(brokerId -> { - assertFalse(active.replicationControl().isBrokerUnfenced(brokerId), + assertFalse(active.clusterControl().unfenced(brokerId), "Broker " + brokerId + " should have been fenced"); }); @@ -346,7 +346,7 @@ public void testBalancePartitionLeaders() throws Throwable { () -> { sendBrokerheartbeat(active, brokersToKeepUnfenced, brokerEpochs); for (Integer brokerId : brokersToFence) { - if (active.replicationControl().isBrokerUnfenced(brokerId)) { + if (active.clusterControl().unfenced(brokerId)) { return false; } } @@ -361,11 +361,11 @@ public void testBalancePartitionLeaders() throws Throwable { // At this point only the brokers we want fenced should be fenced. brokersToKeepUnfenced.forEach(brokerId -> { - assertTrue(active.replicationControl().isBrokerUnfenced(brokerId), + assertTrue(active.clusterControl().unfenced(brokerId), "Broker " + brokerId + " should have been unfenced"); }); brokersToFence.forEach(brokerId -> { - assertFalse(active.replicationControl().isBrokerUnfenced(brokerId), + assertFalse(active.clusterControl().unfenced(brokerId), "Broker " + brokerId + " should have been fenced"); }); From c3f8559a38426f83d668b0f45f1dc506f2fe103c Mon Sep 17 00:00:00 2001 From: David Jacot Date: Tue, 7 Jun 2022 11:46:20 +0200 Subject: [PATCH 18/18] cleanup --- .../org/apache/kafka/controller/QuorumController.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) 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 bf4632d3629fe..0c6d665027e3f 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -1563,10 +1563,10 @@ private QuorumController(LogContext logContext, build(); this.clientQuotaControlManager = new ClientQuotaControlManager(snapshotRegistry); this.featureControl = new FeatureControlManager.Builder(). - setLogContext(logContext). - setQuorumFeatures(quorumFeatures). - setSnapshotRegistry(snapshotRegistry). - build(); + setLogContext(logContext). + setQuorumFeatures(quorumFeatures). + setSnapshotRegistry(snapshotRegistry). + build(); this.clusterControl = new ClusterControlManager.Builder(). setLogContext(logContext). setClusterId(clusterId).