From a84e493fc643b82ea3d9ee0b7ec9ad5f006a07b9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jos=C3=A9=20Armando=20Garc=C3=ADa=20Sancio?= Date: Wed, 3 Mar 2021 09:35:31 -0800 Subject: [PATCH 1/4] KAFKA-12376: Apply atomic append to the log --- .../main/scala/kafka/raft/RaftManager.scala | 29 ++++++- .../controller/ClientQuotaControlManager.java | 2 +- .../ConfigurationControlManager.java | 4 +- .../kafka/controller/ControllerResult.java | 26 ++++-- .../controller/FeatureControlManager.java | 2 +- .../kafka/controller/QuorumController.java | 13 ++- .../controller/ReplicationControlManager.java | 2 +- .../apache/kafka/metalog/LocalLogManager.java | 17 +++- .../apache/kafka/metalog/MetaLogManager.java | 23 ++++- .../ConfigurationControlManagerTest.java | 86 ++++++++++++++----- .../controller/FeatureControlManagerTest.java | 47 +++++++--- .../apache/kafka/metalog/LocalLogManager.java | 17 +++- .../kafka/raft/metadata/MetaLogRaftShim.java | 28 +++++- 13 files changed, 234 insertions(+), 62 deletions(-) diff --git a/core/src/main/scala/kafka/raft/RaftManager.scala b/core/src/main/scala/kafka/raft/RaftManager.scala index ecf89348ccaf9..2b00b1c350dc9 100644 --- a/core/src/main/scala/kafka/raft/RaftManager.scala +++ b/core/src/main/scala/kafka/raft/RaftManager.scala @@ -91,6 +91,11 @@ trait RaftManager[T] { listener: RaftClient.Listener[T] ): Unit + def scheduleAtomicAppend( + epoch: Int, + records: Seq[T] + ): Option[Long] + def scheduleAppend( epoch: Int, records: Seq[T] @@ -156,16 +161,32 @@ class KafkaRaftManager[T]( raftClient.register(listener) } + override def scheduleAtomicAppend( + epoch: Int, + records: Seq[T] + ): Option[Long] = { + append(epoch, records, true) + } + override def scheduleAppend( epoch: Int, records: Seq[T] ): Option[Long] = { - val offset: java.lang.Long = raftClient.scheduleAppend(epoch, records.asJava) - if (offset == null) { - None + append(epoch, records, false) + } + + private def append( + epoch: Int, + records: Seq[T], + isAtomic: Boolean + ): Option[Long] = { + val offset = if (isAtomic) { + raftClient.scheduleAtomicAppend(epoch, records.asJava) } else { - Some(Long.unbox(offset)) + raftClient.scheduleAppend(epoch, records.asJava) } + + Option(offset).map(Long.unbox) } override def handleRequest( diff --git a/metadata/src/main/java/org/apache/kafka/controller/ClientQuotaControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ClientQuotaControlManager.java index 4aac9e4882f46..2221fcfc5d683 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClientQuotaControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClientQuotaControlManager.java @@ -86,7 +86,7 @@ ControllerResult> alterClientQuotas( } }); - return new ControllerResult<>(outputRecords, outputResults); + return new ControllerResult<>(outputRecords, outputResults, true); } /** diff --git a/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java index 4402b3a117d83..1a5a14ad9412b 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java @@ -83,7 +83,7 @@ ControllerResult> incrementalAlterConfigs( outputRecords, outputResults); } - return new ControllerResult<>(outputRecords, outputResults); + return new ControllerResult<>(outputRecords, outputResults, true); } private void incrementalAlterConfigResource(ConfigResource configResource, @@ -171,7 +171,7 @@ ControllerResult> legacyAlterConfigs( outputRecords, outputResults); } - return new ControllerResult<>(outputRecords, outputResults); + return new ControllerResult<>(outputRecords, outputResults, true); } private void legacyAlterConfigResource(ConfigResource configResource, diff --git a/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java b/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java index 4906c8b0972a8..464f08f78fe6f 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java @@ -28,15 +28,21 @@ class ControllerResult { private final List records; private final T response; + private final boolean isAtomic; public ControllerResult(T response) { this(new ArrayList<>(), response); } public ControllerResult(List records, T response) { + this(records, response, false); + } + + public ControllerResult(List records, T response, boolean isAtomic) { Objects.requireNonNull(records); this.records = records; this.response = response; + this.isAtomic = isAtomic; } public List records() { @@ -47,6 +53,10 @@ public T response() { return response; } + public boolean isAtomic() { + return isAtomic; + } + @Override public boolean equals(Object o) { if (o == null || (!o.getClass().equals(getClass()))) { @@ -54,22 +64,26 @@ public boolean equals(Object o) { } ControllerResult other = (ControllerResult) o; return records.equals(other.records) && - Objects.equals(response, other.response); + Objects.equals(response, other.response) && + Objects.equals(isAtomic, other.isAtomic); } @Override public int hashCode() { - return Objects.hash(records, response); + return Objects.hash(records, response, isAtomic); } @Override public String toString() { - return "ControllerResult(records=" + String.join(",", - records.stream().map(r -> r.toString()).collect(Collectors.toList())) + - ", response=" + response + ")"; + return String.format( + "ControllerResult(records=%s, response=%s, isAtomic=%s)", + String.join(",", records.stream().map(r -> r.toString()).collect(Collectors.toList())), + response, + isAtomic + ); } public ControllerResult withoutRecords() { - return new ControllerResult<>(new ArrayList<>(), response); + return new ControllerResult<>(new ArrayList<>(), response, false); } } diff --git a/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java index 25ff3fdcdd80e..ce5e5505659f7 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java @@ -69,7 +69,7 @@ ControllerResult> updateFeatures( results.put(entry.getKey(), updateFeature(entry.getKey(), entry.getValue(), downgradeables.contains(entry.getKey()), brokerFeatures, records)); } - return new ControllerResult<>(records, results); + return new ControllerResult<>(records, results, records.isEmpty() ? false : true); } private ApiError updateFeature(String featureName, 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 198097538985a..a0958781bd9d8 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -265,7 +265,7 @@ private Throwable handleEventException(String name, class ControlEvent implements EventQueue.Event { private final String name; private final Runnable handler; - private long eventCreatedTimeNs = time.nanoseconds(); + private final long eventCreatedTimeNs = time.nanoseconds(); private Optional startProcessingTimeNs = Optional.empty(); ControlEvent(String name, Runnable handler) { @@ -307,7 +307,7 @@ class ControllerReadEvent implements EventQueue.Event { private final String name; private final CompletableFuture future; private final Supplier handler; - private long eventCreatedTimeNs = time.nanoseconds(); + private final long eventCreatedTimeNs = time.nanoseconds(); private Optional startProcessingTimeNs = Optional.empty(); ControllerReadEvent(String name, Supplier handler) { @@ -389,7 +389,7 @@ class ControllerWriteEvent implements EventQueue.Event, DeferredEvent { private final String name; private final CompletableFuture future; private final ControllerWriteOperation op; - private long eventCreatedTimeNs = time.nanoseconds(); + private final long eventCreatedTimeNs = time.nanoseconds(); private Optional startProcessingTimeNs = Optional.empty(); private ControllerResultAndOffset resultAndOffset; @@ -441,7 +441,12 @@ public void run() throws Exception { // written before we can return our result to the user. Here, we hand off // the batch of records to the metadata log manager. They will be written // out asynchronously. - long offset = logManager.scheduleWrite(controllerEpoch, result.records()); + final long offset; + if (result.isAtomic()) { + offset = logManager.scheduleAtomicWrite(controllerEpoch, result.records()); + } else { + offset = logManager.scheduleWrite(controllerEpoch, result.records()); + } op.processBatchEndOffset(offset); writeOffset = offset; resultAndOffset = new ControllerResultAndOffset<>(offset, 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 9fd172f486ae6..8f59e38d3e1b4 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java @@ -407,7 +407,7 @@ public void replay(PartitionChangeRecord record) { resultsPrefix = ", "; } log.info("createTopics result(s): {}", resultsBuilder.toString()); - return new ControllerResult<>(records, data); + return new ControllerResult<>(records, data, true); } private ApiError createTopic(CreatableTopic topic, diff --git a/metadata/src/main/java/org/apache/kafka/metalog/LocalLogManager.java b/metadata/src/main/java/org/apache/kafka/metalog/LocalLogManager.java index ef85314e0ef2f..976825275c26d 100644 --- a/metadata/src/main/java/org/apache/kafka/metalog/LocalLogManager.java +++ b/metadata/src/main/java/org/apache/kafka/metalog/LocalLogManager.java @@ -328,8 +328,21 @@ public void register(MetaLogListener listener) throws Exception { @Override public long scheduleWrite(long epoch, List batch) { - return shared.tryAppend(nodeId, leader.epoch(), new LocalRecordBatch( - batch.stream().map(r -> r.message()).collect(Collectors.toList()))); + return scheduleAtomicWrite(epoch, batch); + } + + @Override + public long scheduleAtomicWrite(long epoch, List batch) { + return shared.tryAppend( + nodeId, + leader.epoch(), + new LocalRecordBatch( + batch + .stream() + .map(r -> r.message()) + .collect(Collectors.toList()) + ) + ); } @Override diff --git a/metadata/src/main/java/org/apache/kafka/metalog/MetaLogManager.java b/metadata/src/main/java/org/apache/kafka/metalog/MetaLogManager.java index 67a6ca5385f75..9126245ef3855 100644 --- a/metadata/src/main/java/org/apache/kafka/metalog/MetaLogManager.java +++ b/metadata/src/main/java/org/apache/kafka/metalog/MetaLogManager.java @@ -50,13 +50,30 @@ public interface MetaLogManager { * offset before renouncing its leadership. The listener should determine this by * monitoring the committed offsets. * - * @param epoch The controller epoch. - * @param batch The batch of messages to write. + * @param epoch the controller epoch + * @param batch the batch of messages to write * - * @return The offset of the message. + * @return the offset of the last message in the batch + * @throws IllegalArgumentException if buffer allocatio failed and the client should backoff */ long scheduleWrite(long epoch, List batch); + /** + * Schedule a atomic write to the log. + * + * The write will be scheduled to happen at some time in the future. All of the messages in batch + * will be appended atomically in one batch. The listener may regard the write as successful + * if and only if the MetaLogManager reaches the given offset before renouncing its leadership. + * The listener should determine this by monitoring the committed offsets. + * + * @param epoch the controller epoch + * @param batch the batch of messages to write + * + * @return the offset of the last message in the batch + * @throws IllegalArgumentException if buffer allocatio failed and the client should backoff + */ + long scheduleAtomicWrite(long epoch, List batch); + /** * Renounce the leadership. * diff --git a/metadata/src/test/java/org/apache/kafka/controller/ConfigurationControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/ConfigurationControlManagerTest.java index 49a55338309b5..083ce700bb75b 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ConfigurationControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ConfigurationControlManagerTest.java @@ -135,18 +135,43 @@ public void testIncrementalAlterConfigs() { SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); ConfigurationControlManager manager = new ConfigurationControlManager(new LogContext(), snapshotRegistry, CONFIGS); - assertEquals(new ControllerResult>(Collections.singletonList( - new ApiMessageAndVersion(new ConfigRecord(). - setResourceType(TOPIC.id()).setResourceName("mytopic"). - setName("abc").setValue("123"), (short) 0)), - toMap(entry(BROKER0, new ApiError( - Errors.INVALID_REQUEST, "A DELETE op was given with a non-null value.")), - entry(MYTOPIC, ApiError.NONE))), - manager.incrementalAlterConfigs(toMap(entry(BROKER0, toMap( - entry("foo.bar", entry(DELETE, "abc")), - entry("quux", entry(SET, "abc")))), - entry(MYTOPIC, toMap( - entry("abc", entry(APPEND, "123"))))))); + assertEquals( + new ControllerResult>( + Collections.singletonList( + new ApiMessageAndVersion( + new ConfigRecord() + .setResourceType(TOPIC.id()) + .setResourceName("mytopic") + .setName("abc") + .setValue("123"), + (short) 0 + ) + ), + toMap( + entry( + BROKER0, + new ApiError( + Errors.INVALID_REQUEST, + "A DELETE op was given with a non-null value." + ) + ), + entry(MYTOPIC, ApiError.NONE) + ), + true + ), + manager.incrementalAlterConfigs( + toMap( + entry( + BROKER0, + toMap( + entry("foo.bar", entry(DELETE, "abc")), + entry("quux", entry(SET, "abc")) + ) + ), + entry(MYTOPIC, toMap(entry("abc", entry(APPEND, "123")))) + ) + ) + ); } @Test @@ -184,20 +209,35 @@ public void testLegacyAlterConfigs() { new ApiMessageAndVersion(new ConfigRecord(). setResourceType(TOPIC.id()).setResourceName("mytopic"). setName("def").setValue("901"), (short) 0)); - assertEquals(new ControllerResult>( + assertEquals( + new ControllerResult>( expectedRecords1, - toMap(entry(MYTOPIC, ApiError.NONE))), - manager.legacyAlterConfigs(toMap(entry(MYTOPIC, toMap( - entry("abc", "456"), entry("def", "901")))))); + toMap(entry(MYTOPIC, ApiError.NONE)), + true + ), + manager.legacyAlterConfigs( + toMap(entry(MYTOPIC, toMap(entry("abc", "456"), entry("def", "901")))) + ) + ); for (ApiMessageAndVersion message : expectedRecords1) { manager.replay((ConfigRecord) message.message()); } - assertEquals(new ControllerResult>(Arrays.asList( - new ApiMessageAndVersion(new ConfigRecord(). - setResourceType(TOPIC.id()).setResourceName("mytopic"). - setName("abc").setValue(null), (short) 0)), - toMap(entry(MYTOPIC, ApiError.NONE))), - manager.legacyAlterConfigs(toMap(entry(MYTOPIC, toMap( - entry("def", "901")))))); + assertEquals( + new ControllerResult>( + Arrays.asList( + new ApiMessageAndVersion( + new ConfigRecord() + .setResourceType(TOPIC.id()) + .setResourceName("mytopic") + .setName("abc") + .setValue(null), + (short) 0 + ) + ), + toMap(entry(MYTOPIC, ApiError.NONE)), + true + ), + manager.legacyAlterConfigs(toMap(entry(MYTOPIC, toMap(entry("def", "901"))))) + ); } } diff --git a/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java index 8687cc8f562d8..c34954ff138c7 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java @@ -101,12 +101,23 @@ public void testUpdateFeaturesErrorCases() { SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); FeatureControlManager manager = new FeatureControlManager( rangeMap("foo", 1, 5, "bar", 1, 2), snapshotRegistry); - assertEquals(new ControllerResult<>(Collections. - singletonMap("foo", new ApiError(Errors.INVALID_UPDATE_VERSION, - "Broker 5 does not support the given feature range."))), - manager.updateFeatures(rangeMap("foo", 1, 3), + + assertEquals( + new ControllerResult<>( + Collections.singletonMap( + "foo", + new ApiError( + Errors.INVALID_UPDATE_VERSION, + "Broker 5 does not support the given feature range." + ) + ) + ), + manager.updateFeatures( + rangeMap("foo", 1, 3), new HashSet<>(Arrays.asList("foo")), - Collections.singletonMap(5, rangeMap()))); + Collections.singletonMap(5, rangeMap()) + ) + ); ControllerResult> result = manager.updateFeatures( rangeMap("foo", 1, 3), Collections.emptySet(), Collections.emptyMap()); @@ -121,12 +132,24 @@ public void testUpdateFeaturesErrorCases() { manager.updateFeatures(rangeMap("foo", 1, 2), Collections.emptySet(), Collections.emptyMap())); - assertEquals(new ControllerResult<>( - Collections.singletonList(new ApiMessageAndVersion(new FeatureLevelRecord(). - setName("foo").setMinFeatureLevel((short) 1).setMaxFeatureLevel((short) 2), - (short) 0)), - Collections.singletonMap("foo", ApiError.NONE)), - manager.updateFeatures(rangeMap("foo", 1, 2), - new HashSet<>(Collections.singletonList("foo")), Collections.emptyMap())); + assertEquals( + new ControllerResult<>( + Collections.singletonList( + new ApiMessageAndVersion( + new FeatureLevelRecord() + .setName("foo") + .setMinFeatureLevel((short) 1) + .setMaxFeatureLevel((short) 2), + (short) 0 + ) + ), + Collections.singletonMap("foo", ApiError.NONE), + true + ), + manager.updateFeatures( + rangeMap("foo", 1, 2), + new HashSet<>(Collections.singletonList("foo")), Collections.emptyMap() + ) + ); } } diff --git a/metadata/src/test/java/org/apache/kafka/metalog/LocalLogManager.java b/metadata/src/test/java/org/apache/kafka/metalog/LocalLogManager.java index 7b6cf06212e8e..b5f7b38ccc916 100644 --- a/metadata/src/test/java/org/apache/kafka/metalog/LocalLogManager.java +++ b/metadata/src/test/java/org/apache/kafka/metalog/LocalLogManager.java @@ -371,8 +371,21 @@ public void register(MetaLogListener listener) throws Exception { @Override public long scheduleWrite(long epoch, List batch) { - return shared.tryAppend(nodeId, leader.epoch(), new LocalRecordBatch( - batch.stream().map(r -> r.message()).collect(Collectors.toList()))); + return scheduleAtomicWrite(epoch, batch); + } + + @Override + public long scheduleAtomicWrite(long epoch, List batch) { + return shared.tryAppend( + nodeId, + leader.epoch(), + new LocalRecordBatch( + batch + .stream() + .map(r -> r.message()) + .collect(Collectors.toList()) + ) + ); } @Override diff --git a/raft/src/main/java/org/apache/kafka/raft/metadata/MetaLogRaftShim.java b/raft/src/main/java/org/apache/kafka/raft/metadata/MetaLogRaftShim.java index bf88e7d8120a1..1ca63f1b9c3cd 100644 --- a/raft/src/main/java/org/apache/kafka/raft/metadata/MetaLogRaftShim.java +++ b/raft/src/main/java/org/apache/kafka/raft/metadata/MetaLogRaftShim.java @@ -52,9 +52,35 @@ public void register(MetaLogListener listener) { client.register(new ListenerShim(listener)); } + @Override + public long scheduleAtomicWrite(long epoch, List batch) { + return write(epoch, batch, true); + } + @Override public long scheduleWrite(long epoch, List batch) { - return client.scheduleAppend((int) epoch, batch); + return write(epoch, batch, false); + } + + private long write(long epoch, List batch, boolean isAtomic) { + final Long result; + if (isAtomic) { + result = client.scheduleAtomicAppend((int) epoch, batch); + } else { + result = client.scheduleAppend((int) epoch, batch); + } + + if (result == null) { + throw new IllegalArgumentException( + String.format( + "Unable to alloate a buffer for the schedule write operation: epoch %s, batch %s)", + epoch, + batch + ) + ); + } else { + return result; + } } @Override From afad436d54aa8c9761d3e339b47899b3e3d9ee48 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jos=C3=A9=20Armando=20Garc=C3=ADa=20Sancio?= Date: Wed, 3 Mar 2021 13:46:55 -0800 Subject: [PATCH 2/4] Add static construction methods --- .../controller/ClusterControlManager.java | 2 +- .../kafka/controller/ControllerResult.java | 16 ++++----- .../controller/ControllerResultAndOffset.java | 34 +++++++++---------- .../kafka/controller/QuorumController.java | 10 ++---- .../controller/ReplicationControlManager.java | 12 +++---- .../apache/kafka/metalog/LocalLogManager.java | 2 +- .../controller/FeatureControlManagerTest.java | 12 +++---- .../apache/kafka/metalog/LocalLogManager.java | 2 +- 8 files changed, 42 insertions(+), 48 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 6e329c72a0e3f..4748d195986ab 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -213,7 +213,7 @@ public ControllerResult registerBroker( List records = new ArrayList<>(); records.add(new ApiMessageAndVersion(record, (short) 0)); - return new ControllerResult<>(records, new BrokerRegistrationReply(brokerEpoch)); + return ControllerResult.of(records, new BrokerRegistrationReply(brokerEpoch)); } public void replay(RegisterBrokerRecord record) { diff --git a/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java b/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java index 464f08f78fe6f..7402c9dd79d47 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java @@ -30,14 +30,6 @@ class ControllerResult { private final T response; private final boolean isAtomic; - public ControllerResult(T response) { - this(new ArrayList<>(), response); - } - - public ControllerResult(List records, T response) { - this(records, response, false); - } - public ControllerResult(List records, T response, boolean isAtomic) { Objects.requireNonNull(records); this.records = records; @@ -86,4 +78,12 @@ public String toString() { public ControllerResult withoutRecords() { return new ControllerResult<>(new ArrayList<>(), response, false); } + + public static ControllerResult atomicOf(List records, T response) { + return new ControllerResult<>(records, response, true); + } + + public static ControllerResult of(List records, T response) { + return new ControllerResult<>(records, response, false); + } } diff --git a/metadata/src/main/java/org/apache/kafka/controller/ControllerResultAndOffset.java b/metadata/src/main/java/org/apache/kafka/controller/ControllerResultAndOffset.java index 5e483f773d5e2..b888541146e05 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ControllerResultAndOffset.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ControllerResultAndOffset.java @@ -17,26 +17,15 @@ package org.apache.kafka.controller; -import org.apache.kafka.metadata.ApiMessageAndVersion; - -import java.util.ArrayList; -import java.util.List; import java.util.Objects; import java.util.stream.Collectors; -class ControllerResultAndOffset extends ControllerResult { +final class ControllerResultAndOffset extends ControllerResult { private final long offset; - public ControllerResultAndOffset(T response) { - super(new ArrayList<>(), response); - this.offset = -1; - } - - public ControllerResultAndOffset(long offset, - List records, - T response) { - super(records, response); + private ControllerResultAndOffset(long offset, ControllerResult result) { + super(result.records(), result.response(), result.isAtomic()); this.offset = offset; } @@ -52,18 +41,27 @@ public boolean equals(Object o) { ControllerResultAndOffset other = (ControllerResultAndOffset) o; return records().equals(other.records()) && response().equals(other.response()) && + isAtomic() == other.isAtomic() && offset == other.offset; } @Override public int hashCode() { - return Objects.hash(records(), response(), offset); + return Objects.hash(records(), response(), isAtomic(), offset); } @Override public String toString() { - return "ControllerResultAndOffset(records=" + String.join(",", - records().stream().map(r -> r.toString()).collect(Collectors.toList())) + - ", response=" + response() + ", offset=" + offset + ")"; + return String.format( + "ControllerResultAndOffset(records=%s, response=%s, isAtomic=%s, offset=%s)", + String.join(",", records().stream().map(r -> r.toString()).collect(Collectors.toList())), + response(), + isAtomic(), + offset + ); + } + + public static ControllerResultAndOffset of(long offset, ControllerResult result) { + return new ControllerResultAndOffset<>(offset, result); } } 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 a0958781bd9d8..759db1efac790 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -17,7 +17,6 @@ package org.apache.kafka.controller; -import java.util.ArrayList; import java.util.Collection; import java.util.Collections; import java.util.List; @@ -423,8 +422,7 @@ public void run() throws Exception { if (!maybeOffset.isPresent()) { // If the purgatory is empty, there are no pending operations and no // uncommitted state. We can return immediately. - resultAndOffset = new ControllerResultAndOffset<>(-1, - new ArrayList<>(), result.response()); + resultAndOffset = ControllerResultAndOffset.of(-1, result); log.debug("Completing read-only operation {} immediately because " + "the purgatory is empty.", this); complete(null); @@ -432,8 +430,7 @@ public void run() throws Exception { } // If there are operations in the purgatory, we want to wait for the latest // one to complete before returning our result to the user. - resultAndOffset = new ControllerResultAndOffset<>(maybeOffset.get(), - result.records(), result.response()); + resultAndOffset = ControllerResultAndOffset.of(maybeOffset.get(), result); log.debug("Read-only operation {} will be completed when the log " + "reaches offset {}", this, resultAndOffset.offset()); } else { @@ -449,8 +446,7 @@ public void run() throws Exception { } op.processBatchEndOffset(offset); writeOffset = offset; - resultAndOffset = new ControllerResultAndOffset<>(offset, - result.records(), result.response()); + resultAndOffset = ControllerResultAndOffset.of(offset, result); for (ApiMessageAndVersion message : result.records()) { replay(message.message()); } 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 8f59e38d3e1b4..ca571058265b7 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java @@ -407,7 +407,7 @@ public void replay(PartitionChangeRecord record) { resultsPrefix = ", "; } log.info("createTopics result(s): {}", resultsBuilder.toString()); - return new ControllerResult<>(records, data, true); + return ControllerResult.atomicOf(records, data); } private ApiError createTopic(CreatableTopic topic, @@ -626,7 +626,7 @@ ControllerResult alterIsr(AlterIsrRequestData request) { setIsr(partitionData.newIsr())); } } - return new ControllerResult<>(records, response); + return ControllerResult.of(records, response); } /** @@ -780,7 +780,7 @@ ControllerResult electLeaders(ElectLeadersRequestData setErrorMessage(error.message())); } } - return new ControllerResult<>(records, response); + return ControllerResult.of(records, response); } static boolean electionIsUnclean(byte electionType) { @@ -875,7 +875,7 @@ ControllerResult processBrokerHeartbeat( states.next().fenced(), states.next().inControlledShutdown(), states.next().shouldShutDown()); - return new ControllerResult<>(records, reply); + return ControllerResult.of(records, reply); } int bestLeader(int[] replicas, int[] isr, boolean unclean) { @@ -904,7 +904,7 @@ public ControllerResult unregisterBroker(int brokerId) { } List records = new ArrayList<>(); handleBrokerUnregistered(brokerId, registration.epoch(), records); - return new ControllerResult<>(records, null); + return ControllerResult.of(records, null); } ControllerResult maybeFenceStaleBrokers() { @@ -916,6 +916,6 @@ ControllerResult maybeFenceStaleBrokers() { handleBrokerFenced(brokerId, records); heartbeatManager.fence(brokerId); } - return new ControllerResult<>(records, null); + return ControllerResult.of(records, null); } } diff --git a/metadata/src/main/java/org/apache/kafka/metalog/LocalLogManager.java b/metadata/src/main/java/org/apache/kafka/metalog/LocalLogManager.java index 976825275c26d..99ae3a7e9baa8 100644 --- a/metadata/src/main/java/org/apache/kafka/metalog/LocalLogManager.java +++ b/metadata/src/main/java/org/apache/kafka/metalog/LocalLogManager.java @@ -339,7 +339,7 @@ public long scheduleAtomicWrite(long epoch, List batch) { new LocalRecordBatch( batch .stream() - .map(r -> r.message()) + .map(ApiMessageAndVersion::message) .collect(Collectors.toList()) ) ); diff --git a/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java index c34954ff138c7..dde241025925b 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java @@ -61,7 +61,7 @@ public void testUpdateFeatures() { rangeMap("foo", 1, 2), snapshotRegistry); assertEquals(new FeatureMapAndEpoch(new FeatureMap(Collections.emptyMap()), -1), manager.finalizedFeatures(-1)); - assertEquals(new ControllerResult<>(Collections. + assertEquals(ControllerResult.of(Collections.emptyList(), Collections. singletonMap("foo", new ApiError(Errors.INVALID_UPDATE_VERSION, "The controller does not support the given feature range."))), manager.updateFeatures(rangeMap("foo", 1, 3), @@ -103,7 +103,8 @@ public void testUpdateFeaturesErrorCases() { rangeMap("foo", 1, 5, "bar", 1, 2), snapshotRegistry); assertEquals( - new ControllerResult<>( + ControllerResult.of( + Collections.emptyList(), Collections.singletonMap( "foo", new ApiError( @@ -125,7 +126,7 @@ public void testUpdateFeaturesErrorCases() { manager.replay((FeatureLevelRecord) result.records().get(0).message(), 3); snapshotRegistry.createSnapshot(3); - assertEquals(new ControllerResult<>(Collections. + assertEquals(ControllerResult.of(Collections.emptyList(), Collections. singletonMap("foo", new ApiError(Errors.INVALID_UPDATE_VERSION, "Can't downgrade the maximum version of this feature without " + "setting downgradable to true."))), @@ -133,7 +134,7 @@ public void testUpdateFeaturesErrorCases() { Collections.emptySet(), Collections.emptyMap())); assertEquals( - new ControllerResult<>( + ControllerResult.atomicOf( Collections.singletonList( new ApiMessageAndVersion( new FeatureLevelRecord() @@ -143,8 +144,7 @@ public void testUpdateFeaturesErrorCases() { (short) 0 ) ), - Collections.singletonMap("foo", ApiError.NONE), - true + Collections.singletonMap("foo", ApiError.NONE) ), manager.updateFeatures( rangeMap("foo", 1, 2), diff --git a/metadata/src/test/java/org/apache/kafka/metalog/LocalLogManager.java b/metadata/src/test/java/org/apache/kafka/metalog/LocalLogManager.java index b5f7b38ccc916..590f89c3391b3 100644 --- a/metadata/src/test/java/org/apache/kafka/metalog/LocalLogManager.java +++ b/metadata/src/test/java/org/apache/kafka/metalog/LocalLogManager.java @@ -382,7 +382,7 @@ public long scheduleAtomicWrite(long epoch, List batch) { new LocalRecordBatch( batch .stream() - .map(r -> r.message()) + .map(ApiMessageAndVersion::message) .collect(Collectors.toList()) ) ); From de9ba57130f2393c8bd831c8a4c51fe051d86294 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jos=C3=A9=20Armando=20Garc=C3=ADa=20Sancio?= Date: Wed, 3 Mar 2021 13:55:40 -0800 Subject: [PATCH 3/4] Make constructor protected --- .../controller/ClientQuotaControlManager.java | 2 +- .../controller/ConfigurationControlManager.java | 4 ++-- .../apache/kafka/controller/ControllerResult.java | 2 +- .../kafka/controller/FeatureControlManager.java | 7 ++++++- .../ConfigurationControlManagerTest.java | 15 ++++++--------- 5 files changed, 16 insertions(+), 14 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/controller/ClientQuotaControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ClientQuotaControlManager.java index 2221fcfc5d683..9b8e2d683b650 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClientQuotaControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClientQuotaControlManager.java @@ -86,7 +86,7 @@ ControllerResult> alterClientQuotas( } }); - return new ControllerResult<>(outputRecords, outputResults, true); + return ControllerResult.atomicOf(outputRecords, outputResults); } /** diff --git a/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java index 1a5a14ad9412b..5bc82ecafa477 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java @@ -83,7 +83,7 @@ ControllerResult> incrementalAlterConfigs( outputRecords, outputResults); } - return new ControllerResult<>(outputRecords, outputResults, true); + return ControllerResult.atomicOf(outputRecords, outputResults); } private void incrementalAlterConfigResource(ConfigResource configResource, @@ -171,7 +171,7 @@ ControllerResult> legacyAlterConfigs( outputRecords, outputResults); } - return new ControllerResult<>(outputRecords, outputResults, true); + return ControllerResult.atomicOf(outputRecords, outputResults); } private void legacyAlterConfigResource(ConfigResource configResource, diff --git a/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java b/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java index 7402c9dd79d47..5774bf7e0edf8 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java @@ -30,7 +30,7 @@ class ControllerResult { private final T response; private final boolean isAtomic; - public ControllerResult(List records, T response, boolean isAtomic) { + protected ControllerResult(List records, T response, boolean isAtomic) { Objects.requireNonNull(records); this.records = records; this.response = response; diff --git a/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java index ce5e5505659f7..fc540f06373c0 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java @@ -69,7 +69,12 @@ ControllerResult> updateFeatures( results.put(entry.getKey(), updateFeature(entry.getKey(), entry.getValue(), downgradeables.contains(entry.getKey()), brokerFeatures, records)); } - return new ControllerResult<>(records, results, records.isEmpty() ? false : true); + + if (records.isEmpty()) { + return ControllerResult.of(records, results); + } else { + return ControllerResult.atomicOf(records, results); + } } private ApiError updateFeature(String featureName, diff --git a/metadata/src/test/java/org/apache/kafka/controller/ConfigurationControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/ConfigurationControlManagerTest.java index 083ce700bb75b..561a25b2d0ee2 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ConfigurationControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ConfigurationControlManagerTest.java @@ -136,7 +136,7 @@ public void testIncrementalAlterConfigs() { ConfigurationControlManager manager = new ConfigurationControlManager(new LogContext(), snapshotRegistry, CONFIGS); assertEquals( - new ControllerResult>( + ControllerResult.atomicOf( Collections.singletonList( new ApiMessageAndVersion( new ConfigRecord() @@ -156,8 +156,7 @@ public void testIncrementalAlterConfigs() { ) ), entry(MYTOPIC, ApiError.NONE) - ), - true + ) ), manager.incrementalAlterConfigs( toMap( @@ -210,10 +209,9 @@ public void testLegacyAlterConfigs() { setResourceType(TOPIC.id()).setResourceName("mytopic"). setName("def").setValue("901"), (short) 0)); assertEquals( - new ControllerResult>( + ControllerResult.atomicOf( expectedRecords1, - toMap(entry(MYTOPIC, ApiError.NONE)), - true + toMap(entry(MYTOPIC, ApiError.NONE)) ), manager.legacyAlterConfigs( toMap(entry(MYTOPIC, toMap(entry("abc", "456"), entry("def", "901")))) @@ -223,7 +221,7 @@ public void testLegacyAlterConfigs() { manager.replay((ConfigRecord) message.message()); } assertEquals( - new ControllerResult>( + ControllerResult.atomicOf( Arrays.asList( new ApiMessageAndVersion( new ConfigRecord() @@ -234,8 +232,7 @@ public void testLegacyAlterConfigs() { (short) 0 ) ), - toMap(entry(MYTOPIC, ApiError.NONE)), - true + toMap(entry(MYTOPIC, ApiError.NONE)) ), manager.legacyAlterConfigs(toMap(entry(MYTOPIC, toMap(entry("def", "901"))))) ); From 9f5940e3a5992b7c205e84b0b74930ca24057096 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jos=C3=A9=20Armando=20Garc=C3=ADa=20Sancio?= Date: Thu, 4 Mar 2021 09:19:23 -0800 Subject: [PATCH 4/4] Always create an atomic ControllerResult for feature changes --- .../apache/kafka/controller/ControllerResult.java | 6 +++--- .../controller/ControllerResultAndOffset.java | 4 +++- .../kafka/controller/FeatureControlManager.java | 6 +----- .../controller/FeatureControlManagerTest.java | 15 +++++++-------- 4 files changed, 14 insertions(+), 17 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java b/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java index 5774bf7e0edf8..e6ae031b9b3b8 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ControllerResult.java @@ -19,7 +19,7 @@ import org.apache.kafka.metadata.ApiMessageAndVersion; -import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Objects; import java.util.stream.Collectors; @@ -69,14 +69,14 @@ public int hashCode() { public String toString() { return String.format( "ControllerResult(records=%s, response=%s, isAtomic=%s)", - String.join(",", records.stream().map(r -> r.toString()).collect(Collectors.toList())), + String.join(",", records.stream().map(ApiMessageAndVersion::toString).collect(Collectors.toList())), response, isAtomic ); } public ControllerResult withoutRecords() { - return new ControllerResult<>(new ArrayList<>(), response, false); + return new ControllerResult<>(Collections.emptyList(), response, false); } public static ControllerResult atomicOf(List records, T response) { diff --git a/metadata/src/main/java/org/apache/kafka/controller/ControllerResultAndOffset.java b/metadata/src/main/java/org/apache/kafka/controller/ControllerResultAndOffset.java index b888541146e05..8b8ca8dea80da 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ControllerResultAndOffset.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ControllerResultAndOffset.java @@ -17,6 +17,8 @@ package org.apache.kafka.controller; +import org.apache.kafka.metadata.ApiMessageAndVersion; + import java.util.Objects; import java.util.stream.Collectors; @@ -54,7 +56,7 @@ public int hashCode() { public String toString() { return String.format( "ControllerResultAndOffset(records=%s, response=%s, isAtomic=%s, offset=%s)", - String.join(",", records().stream().map(r -> r.toString()).collect(Collectors.toList())), + String.join(",", records().stream().map(ApiMessageAndVersion::toString).collect(Collectors.toList())), response(), isAtomic(), offset diff --git a/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java index fc540f06373c0..99874ac3c5ef7 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java @@ -70,11 +70,7 @@ ControllerResult> updateFeatures( downgradeables.contains(entry.getKey()), brokerFeatures, records)); } - if (records.isEmpty()) { - return ControllerResult.of(records, results); - } else { - return ControllerResult.atomicOf(records, results); - } + return ControllerResult.atomicOf(records, results); } private ApiError updateFeature(String featureName, diff --git a/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java index dde241025925b..0670984e52876 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java @@ -18,10 +18,8 @@ package org.apache.kafka.controller; import java.util.ArrayList; -import java.util.Arrays; import java.util.Collections; import java.util.HashMap; -import java.util.HashSet; import java.util.List; import java.util.Map; import org.apache.kafka.common.metadata.FeatureLevelRecord; @@ -61,11 +59,11 @@ public void testUpdateFeatures() { rangeMap("foo", 1, 2), snapshotRegistry); assertEquals(new FeatureMapAndEpoch(new FeatureMap(Collections.emptyMap()), -1), manager.finalizedFeatures(-1)); - assertEquals(ControllerResult.of(Collections.emptyList(), Collections. + assertEquals(ControllerResult.atomicOf(Collections.emptyList(), Collections. singletonMap("foo", new ApiError(Errors.INVALID_UPDATE_VERSION, "The controller does not support the given feature range."))), manager.updateFeatures(rangeMap("foo", 1, 3), - new HashSet<>(Arrays.asList("foo")), + Collections.singleton("foo"), Collections.emptyMap())); ControllerResult> result = manager.updateFeatures( rangeMap("foo", 1, 2, "bar", 1, 1), Collections.emptySet(), @@ -103,7 +101,7 @@ public void testUpdateFeaturesErrorCases() { rangeMap("foo", 1, 5, "bar", 1, 2), snapshotRegistry); assertEquals( - ControllerResult.of( + ControllerResult.atomicOf( Collections.emptyList(), Collections.singletonMap( "foo", @@ -115,7 +113,7 @@ public void testUpdateFeaturesErrorCases() { ), manager.updateFeatures( rangeMap("foo", 1, 3), - new HashSet<>(Arrays.asList("foo")), + Collections.singleton("foo"), Collections.singletonMap(5, rangeMap()) ) ); @@ -126,7 +124,7 @@ public void testUpdateFeaturesErrorCases() { manager.replay((FeatureLevelRecord) result.records().get(0).message(), 3); snapshotRegistry.createSnapshot(3); - assertEquals(ControllerResult.of(Collections.emptyList(), Collections. + assertEquals(ControllerResult.atomicOf(Collections.emptyList(), Collections. singletonMap("foo", new ApiError(Errors.INVALID_UPDATE_VERSION, "Can't downgrade the maximum version of this feature without " + "setting downgradable to true."))), @@ -148,7 +146,8 @@ public void testUpdateFeaturesErrorCases() { ), manager.updateFeatures( rangeMap("foo", 1, 2), - new HashSet<>(Collections.singletonList("foo")), Collections.emptyMap() + Collections.singleton("foo"), + Collections.emptyMap() ) ); }