diff --git a/LICENSE-binary b/LICENSE-binary
index 842962e61ad2c..131d64b863373 100644
--- a/LICENSE-binary
+++ b/LICENSE-binary
@@ -313,7 +313,8 @@ argparse4j-0.7.0, see: licenses/argparse-MIT
jopt-simple-5.0.4, see: licenses/jopt-simple-MIT
slf4j-api-1.7.36, see: licenses/slf4j-MIT
slf4j-reload4j-1.7.36, see: licenses/slf4j-MIT
-classgraph-4.8.138, see: license/classgraph-MIT
+classgraph-4.8.138, see: licenses/classgraph-MIT
+pcollections-4.0.1, see: licenses/pcollections-MIT
---------------------------------------
BSD 2-Clause
diff --git a/build.gradle b/build.gradle
index a4ea31ceed52b..773aa7a776962 100644
--- a/build.gradle
+++ b/build.gradle
@@ -1222,6 +1222,10 @@ project(':metadata') {
javadoc {
enabled = false
}
+
+ checkstyle {
+ configProperties = checkstyleConfigProperties("import-control-metadata.xml")
+ }
}
project(':group-coordinator') {
@@ -1554,11 +1558,13 @@ project(':server-common') {
implementation libs.slf4jApi
implementation libs.metrics
implementation libs.joptSimple
+ implementation libs.pcollections
testImplementation project(':clients')
testImplementation project(':clients').sourceSets.test.output
testImplementation libs.junitJupiter
testImplementation libs.mockitoCore
+ testImplementation libs.mockitoInline // supports mocking static methods, final classes, etc.
testImplementation libs.hamcrest
testRuntimeOnly libs.slf4jlog4j
@@ -1605,6 +1611,10 @@ project(':server-common') {
clean.doFirst {
delete "$buildDir/kafka/"
}
+
+ checkstyle {
+ configProperties = checkstyleConfigProperties("import-control-server-common.xml")
+ }
}
project(':storage:api') {
diff --git a/checkstyle/import-control-jmh-benchmarks.xml b/checkstyle/import-control-jmh-benchmarks.xml
index d6e966498d8a3..4cbd34f89cbfd 100644
--- a/checkstyle/import-control-jmh-benchmarks.xml
+++ b/checkstyle/import-control-jmh-benchmarks.xml
@@ -51,6 +51,7 @@
+
diff --git a/checkstyle/import-control-metadata.xml b/checkstyle/import-control-metadata.xml
new file mode 100644
index 0000000000000..ade866d6a0706
--- /dev/null
+++ b/checkstyle/import-control-metadata.xml
@@ -0,0 +1,176 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/checkstyle/import-control-server-common.xml b/checkstyle/import-control-server-common.xml
new file mode 100644
index 0000000000000..d310d81a8321e
--- /dev/null
+++ b/checkstyle/import-control-server-common.xml
@@ -0,0 +1,82 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/checkstyle/import-control.xml b/checkstyle/import-control.xml
index 72948e541d5c9..c3318807f206f 100644
--- a/checkstyle/import-control.xml
+++ b/checkstyle/import-control.xml
@@ -88,14 +88,6 @@
-
-
-
-
-
-
-
-
@@ -206,122 +198,6 @@
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
-
@@ -358,19 +234,10 @@
-
-
-
-
-
-
-
-
-
diff --git a/core/src/test/scala/unit/kafka/server/metadata/BrokerMetadataPublisherTest.scala b/core/src/test/scala/unit/kafka/server/metadata/BrokerMetadataPublisherTest.scala
index ffe9b3f40764b..d7d332ea7fb28 100644
--- a/core/src/test/scala/unit/kafka/server/metadata/BrokerMetadataPublisherTest.scala
+++ b/core/src/test/scala/unit/kafka/server/metadata/BrokerMetadataPublisherTest.scala
@@ -177,9 +177,9 @@ class BrokerMetadataPublisherTest {
private def topicsImage(
topics: Seq[TopicImage]
): TopicsImage = {
- val idsMap = topics.map(t => t.id -> t).toMap
- val namesMap = topics.map(t => t.name -> t).toMap
- new TopicsImage(idsMap.asJava, namesMap.asJava)
+ var retval = TopicsImage.EMPTY
+ topics.foreach { t => retval = retval.including(t) }
+ retval
}
private def newMockDynamicConfigPublisher(
diff --git a/gradle/dependencies.gradle b/gradle/dependencies.gradle
index b3308b6a1352f..3d3778bae3818 100644
--- a/gradle/dependencies.gradle
+++ b/gradle/dependencies.gradle
@@ -108,6 +108,7 @@ versions += [
metrics: "2.2.0",
mockito: "4.9.0",
netty: "4.1.86.Final",
+ pcollections: "4.0.1",
powermock: "2.0.9",
reflections: "0.9.12",
reload4j: "1.2.19",
@@ -198,6 +199,7 @@ libs += [
mockitoJunitJupiter: "org.mockito:mockito-junit-jupiter:$versions.mockito",
nettyHandler: "io.netty:netty-handler:$versions.netty",
nettyTransportNativeEpoll: "io.netty:netty-transport-native-epoll:$versions.netty",
+ pcollections: "org.pcollections:pcollections:$versions.pcollections",
powermockJunit4: "org.powermock:powermock-module-junit4:$versions.powermock",
powermockEasymock: "org.powermock:powermock-api-easymock:$versions.powermock",
reflections: "org.reflections:reflections:$versions.reflections",
diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/metadata/KRaftMetadataRequestBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/metadata/KRaftMetadataRequestBenchmark.java
new file mode 100644
index 0000000000000..8e997385fc362
--- /dev/null
+++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/metadata/KRaftMetadataRequestBenchmark.java
@@ -0,0 +1,235 @@
+/*
+ * 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.jmh.metadata;
+
+import kafka.coordinator.transaction.TransactionCoordinator;
+import kafka.network.RequestChannel;
+import kafka.network.RequestConvertToJson;
+import kafka.server.AutoTopicCreationManager;
+import kafka.server.BrokerTopicStats;
+import kafka.server.ClientQuotaManager;
+import kafka.server.ClientRequestQuotaManager;
+import kafka.server.ControllerMutationQuotaManager;
+import kafka.server.FetchManager;
+import kafka.server.ForwardingManager;
+import kafka.server.KafkaApis;
+import kafka.server.KafkaConfig;
+import kafka.server.KafkaConfig$;
+import kafka.server.MetadataCache;
+import kafka.server.QuotaFactory;
+import kafka.server.RaftSupport;
+import kafka.server.ReplicaManager;
+import kafka.server.ReplicationQuotaManager;
+import kafka.server.SimpleApiVersionManager;
+import kafka.server.builders.KafkaApisBuilder;
+import kafka.server.metadata.KRaftMetadataCache;
+import kafka.server.metadata.MockConfigRepository;
+import org.apache.kafka.common.Uuid;
+import org.apache.kafka.common.memory.MemoryPool;
+import org.apache.kafka.common.message.ApiMessageType;
+import org.apache.kafka.common.message.UpdateMetadataRequestData.UpdateMetadataEndpoint;
+import org.apache.kafka.common.metadata.PartitionRecord;
+import org.apache.kafka.common.metadata.RegisterBrokerRecord;
+import org.apache.kafka.common.metadata.TopicRecord;
+import org.apache.kafka.common.metrics.Metrics;
+import org.apache.kafka.common.network.ClientInformation;
+import org.apache.kafka.common.network.ListenerName;
+import org.apache.kafka.common.requests.MetadataRequest;
+import org.apache.kafka.common.requests.RequestContext;
+import org.apache.kafka.common.requests.RequestHeader;
+import org.apache.kafka.common.security.auth.KafkaPrincipal;
+import org.apache.kafka.common.security.auth.SecurityProtocol;
+import org.apache.kafka.common.utils.Time;
+import org.apache.kafka.coordinator.group.GroupCoordinator;
+import org.apache.kafka.image.MetadataDelta;
+import org.apache.kafka.image.MetadataImage;
+import org.apache.kafka.image.MetadataProvenance;
+import org.mockito.Mockito;
+import org.openjdk.jmh.annotations.Benchmark;
+import org.openjdk.jmh.annotations.BenchmarkMode;
+import org.openjdk.jmh.annotations.Fork;
+import org.openjdk.jmh.annotations.Level;
+import org.openjdk.jmh.annotations.Measurement;
+import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Param;
+import org.openjdk.jmh.annotations.Scope;
+import org.openjdk.jmh.annotations.Setup;
+import org.openjdk.jmh.annotations.State;
+import org.openjdk.jmh.annotations.TearDown;
+import org.openjdk.jmh.annotations.Warmup;
+import scala.Option;
+
+import java.nio.ByteBuffer;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.Optional;
+import java.util.Properties;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.IntStream;
+
+@State(Scope.Benchmark)
+@Fork(value = 1)
+@Warmup(iterations = 5)
+@Measurement(iterations = 15)
+@BenchmarkMode(Mode.AverageTime)
+@OutputTimeUnit(TimeUnit.NANOSECONDS)
+
+public class KRaftMetadataRequestBenchmark {
+ @Param({"500", "1000", "5000"})
+ private int topicCount;
+ @Param({"10", "20", "50"})
+ private int partitionCount;
+
+ private RequestChannel requestChannel = Mockito.mock(RequestChannel.class, Mockito.withSettings().stubOnly());
+ private RequestChannel.Metrics requestChannelMetrics = Mockito.mock(RequestChannel.Metrics.class);
+ private ReplicaManager replicaManager = Mockito.mock(ReplicaManager.class);
+ private GroupCoordinator groupCoordinator = Mockito.mock(GroupCoordinator.class);
+ private TransactionCoordinator transactionCoordinator = Mockito.mock(TransactionCoordinator.class);
+ private AutoTopicCreationManager autoTopicCreationManager = Mockito.mock(AutoTopicCreationManager.class);
+ private Metrics metrics = new Metrics();
+ private int brokerId = 1;
+ private ForwardingManager forwardingManager = Mockito.mock(ForwardingManager.class);
+ private KRaftMetadataCache metadataCache = MetadataCache.kRaftMetadataCache(brokerId);
+ private ClientQuotaManager clientQuotaManager = Mockito.mock(ClientQuotaManager.class);
+ private ClientRequestQuotaManager clientRequestQuotaManager = Mockito.mock(ClientRequestQuotaManager.class);
+ private ControllerMutationQuotaManager controllerMutationQuotaManager = Mockito.mock(ControllerMutationQuotaManager.class);
+ private ReplicationQuotaManager replicaQuotaManager = Mockito.mock(ReplicationQuotaManager.class);
+ private QuotaFactory.QuotaManagers quotaManagers = new QuotaFactory.QuotaManagers(clientQuotaManager,
+ clientQuotaManager, clientRequestQuotaManager, controllerMutationQuotaManager, replicaQuotaManager,
+ replicaQuotaManager, replicaQuotaManager, Option.empty());
+ private FetchManager fetchManager = Mockito.mock(FetchManager.class);
+ private BrokerTopicStats brokerTopicStats = new BrokerTopicStats();
+ private KafkaPrincipal principal = new KafkaPrincipal(KafkaPrincipal.USER_TYPE, "test-user");
+ private KafkaApis kafkaApis;
+ private RequestChannel.Request allTopicMetadataRequest;
+
+ @Setup(Level.Trial)
+ public void setup() {
+ initializeMetadataCache();
+ kafkaApis = createKafkaApis();
+ allTopicMetadataRequest = buildAllTopicMetadataRequest();
+ }
+
+ private void initializeMetadataCache() {
+ MetadataDelta buildupMetadataDelta = new MetadataDelta(MetadataImage.EMPTY);
+ IntStream.range(0, 5).forEach(brokerId -> {
+ RegisterBrokerRecord.BrokerEndpointCollection endpoints = new RegisterBrokerRecord.BrokerEndpointCollection();
+ endpoints(brokerId).forEach(endpoint ->
+ endpoints.add(new RegisterBrokerRecord.BrokerEndpoint().
+ setHost(endpoint.host()).
+ setPort(endpoint.port()).
+ setName(endpoint.listener()).
+ setSecurityProtocol(endpoint.securityProtocol())));
+ buildupMetadataDelta.replay(new RegisterBrokerRecord().
+ setBrokerId(brokerId).
+ setBrokerEpoch(100L).
+ setFenced(false).
+ setRack(null).
+ setEndPoints(endpoints).
+ setIncarnationId(Uuid.fromString(Uuid.randomUuid().toString())));
+ });
+ IntStream.range(0, topicCount).forEach(topicNum -> {
+ Uuid topicId = Uuid.randomUuid();
+ buildupMetadataDelta.replay(new TopicRecord().setName("topic-" + topicNum).setTopicId(topicId));
+ IntStream.range(0, partitionCount).forEach(partitionId ->
+ buildupMetadataDelta.replay(new PartitionRecord().
+ setPartitionId(partitionId).
+ setTopicId(topicId).
+ setReplicas(Arrays.asList(0, 1, 3)).
+ setIsr(Arrays.asList(0, 1, 3)).
+ setRemovingReplicas(Collections.emptyList()).
+ setAddingReplicas(Collections.emptyList()).
+ setLeader(partitionCount % 5).
+ setLeaderEpoch(0)));
+ });
+ metadataCache.setImage(buildupMetadataDelta.apply(MetadataProvenance.EMPTY));
+ }
+
+ private List endpoints(final int brokerId) {
+ return Collections.singletonList(
+ new UpdateMetadataEndpoint()
+ .setHost("host_" + brokerId)
+ .setPort(9092)
+ .setSecurityProtocol(SecurityProtocol.PLAINTEXT.id)
+ .setListener(ListenerName.forSecurityProtocol(SecurityProtocol.PLAINTEXT).value()));
+ }
+
+ private KafkaApis createKafkaApis() {
+ Properties kafkaProps = new Properties();
+ kafkaProps.put(KafkaConfig$.MODULE$.NodeIdProp(), brokerId + "");
+ kafkaProps.put(KafkaConfig$.MODULE$.ProcessRolesProp(), "broker");
+ kafkaProps.put(KafkaConfig$.MODULE$.QuorumVotersProp(), "9000@foo:8092");
+ kafkaProps.put(KafkaConfig$.MODULE$.ControllerListenerNamesProp(), "CONTROLLER");
+ KafkaConfig config = new KafkaConfig(kafkaProps);
+ return new KafkaApisBuilder().
+ setRequestChannel(requestChannel).
+ setMetadataSupport(new RaftSupport(forwardingManager, metadataCache)).
+ setReplicaManager(replicaManager).
+ setGroupCoordinator(groupCoordinator).
+ setTxnCoordinator(transactionCoordinator).
+ setAutoTopicCreationManager(autoTopicCreationManager).
+ setBrokerId(brokerId).
+ setConfig(config).
+ setConfigRepository(new MockConfigRepository()).
+ setMetadataCache(metadataCache).
+ setMetrics(metrics).
+ setAuthorizer(Optional.empty()).
+ setQuotas(quotaManagers).
+ setFetchManager(fetchManager).
+ setBrokerTopicStats(brokerTopicStats).
+ setClusterId("clusterId").
+ setTime(Time.SYSTEM).
+ setTokenManager(null).
+ setApiVersionManager(new SimpleApiVersionManager(ApiMessageType.ListenerType.BROKER, false)).
+ build();
+ }
+
+ @TearDown(Level.Trial)
+ public void tearDown() {
+ kafkaApis.close();
+ metrics.close();
+ }
+
+ private RequestChannel.Request buildAllTopicMetadataRequest() {
+ MetadataRequest metadataRequest = MetadataRequest.Builder.allTopics().build();
+ RequestHeader header = new RequestHeader(metadataRequest.apiKey(), metadataRequest.version(), "", 0);
+ ByteBuffer bodyBuffer = metadataRequest.serialize();
+
+ RequestContext context = new RequestContext(header, "1", null, principal,
+ ListenerName.forSecurityProtocol(SecurityProtocol.PLAINTEXT),
+ SecurityProtocol.PLAINTEXT, ClientInformation.EMPTY, false);
+ return new RequestChannel.Request(1, context, 0, MemoryPool.NONE, bodyBuffer, requestChannelMetrics, Option.empty());
+ }
+
+ @Benchmark
+ public void testMetadataRequestForAllTopics() {
+ kafkaApis.handleTopicMetadataRequest(allTopicMetadataRequest);
+ }
+
+ @Benchmark
+ public String testRequestToJson() {
+ return RequestConvertToJson.requestDesc(allTopicMetadataRequest.header(), allTopicMetadataRequest.requestLog(), allTopicMetadataRequest.isForwarded()).toString();
+ }
+
+ @Benchmark
+ public void testTopicIdInfo() {
+ metadataCache.topicIdInfo();
+ }
+}
diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/metadata/TopicsImageSingleRecordChangeBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/metadata/TopicsImageSingleRecordChangeBenchmark.java
new file mode 100644
index 0000000000000..8f88a7a1e6f52
--- /dev/null
+++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/metadata/TopicsImageSingleRecordChangeBenchmark.java
@@ -0,0 +1,90 @@
+/*
+ * 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.jmh.metadata;
+
+import org.apache.kafka.common.Uuid;
+import org.apache.kafka.common.metadata.PartitionRecord;
+import org.apache.kafka.common.metadata.TopicRecord;
+import org.apache.kafka.image.TopicsDelta;
+import org.apache.kafka.image.TopicsImage;
+import org.openjdk.jmh.annotations.Benchmark;
+import org.openjdk.jmh.annotations.BenchmarkMode;
+import org.openjdk.jmh.annotations.Fork;
+import org.openjdk.jmh.annotations.Level;
+import org.openjdk.jmh.annotations.Measurement;
+import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Param;
+import org.openjdk.jmh.annotations.Scope;
+import org.openjdk.jmh.annotations.Setup;
+import org.openjdk.jmh.annotations.State;
+import org.openjdk.jmh.annotations.Warmup;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.concurrent.TimeUnit;
+
+@State(Scope.Benchmark)
+@Fork(value = 1)
+@Warmup(iterations = 3)
+@Measurement(iterations = 5)
+@BenchmarkMode(Mode.AverageTime)
+@OutputTimeUnit(TimeUnit.NANOSECONDS)
+public class TopicsImageSingleRecordChangeBenchmark {
+ @Param({"12500", "25000", "50000", "100000"})
+ private int totalTopicCount;
+ @Param({"10"})
+ private int partitionsPerTopic;
+ @Param({"3"})
+ private int replicationFactor;
+ @Param({"10000"})
+ private int numReplicasPerBroker;
+
+ private TopicsDelta topicsDelta;
+
+
+ @Setup(Level.Trial)
+ public void setup() {
+ // build an image containing all the specified topics and partitions
+ TopicsDelta buildupTopicsDelta = TopicsImageSnapshotLoadBenchmark.getInitialTopicsDelta(totalTopicCount, partitionsPerTopic, replicationFactor, numReplicasPerBroker);
+ TopicsImage builtupTopicsImage = buildupTopicsDelta.apply();
+ // build a delta to apply within the benchmark code
+ // that adds a single topic-partition
+ topicsDelta = new TopicsDelta(builtupTopicsImage);
+ Uuid newTopicUuid = Uuid.randomUuid();
+ TopicRecord newTopicRecord = new TopicRecord().setName("newtopic").setTopicId(newTopicUuid);
+ topicsDelta.replay(newTopicRecord);
+ ArrayList replicas = TopicsImageSnapshotLoadBenchmark.getReplicas(totalTopicCount, partitionsPerTopic, replicationFactor, numReplicasPerBroker, 0);
+ ArrayList isr = new ArrayList<>(replicas);
+ PartitionRecord newPartitionRecord = new PartitionRecord().
+ setPartitionId(0).
+ setTopicId(newTopicUuid).
+ setReplicas(replicas).
+ setIsr(isr).
+ setRemovingReplicas(Collections.emptyList()).
+ setAddingReplicas(Collections.emptyList()).
+ setLeader(0);
+ topicsDelta.replay(newPartitionRecord);
+ System.out.print("(Adding a single topic to metadata having " + totalTopicCount + " total topics) ");
+ }
+
+ @Benchmark
+ public void testTopicsDeltaSingleTopicAdd() {
+ topicsDelta.apply();
+ }
+}
diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/metadata/TopicsImageSnapshotLoadBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/metadata/TopicsImageSnapshotLoadBenchmark.java
new file mode 100644
index 0000000000000..10961d9d0252b
--- /dev/null
+++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/metadata/TopicsImageSnapshotLoadBenchmark.java
@@ -0,0 +1,112 @@
+/*
+ * 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.jmh.metadata;
+
+import org.apache.kafka.common.Uuid;
+import org.apache.kafka.common.metadata.PartitionRecord;
+import org.apache.kafka.common.metadata.TopicRecord;
+import org.apache.kafka.image.TopicsDelta;
+import org.apache.kafka.image.TopicsImage;
+import org.openjdk.jmh.annotations.Benchmark;
+import org.openjdk.jmh.annotations.BenchmarkMode;
+import org.openjdk.jmh.annotations.Fork;
+import org.openjdk.jmh.annotations.Level;
+import org.openjdk.jmh.annotations.Measurement;
+import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Param;
+import org.openjdk.jmh.annotations.Scope;
+import org.openjdk.jmh.annotations.Setup;
+import org.openjdk.jmh.annotations.State;
+import org.openjdk.jmh.annotations.Warmup;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.stream.IntStream;
+
+@State(Scope.Benchmark)
+@Fork(value = 1)
+@Warmup(iterations = 3)
+@Measurement(iterations = 5)
+@BenchmarkMode(Mode.AverageTime)
+@OutputTimeUnit(TimeUnit.MILLISECONDS)
+public class TopicsImageSnapshotLoadBenchmark {
+ @Param({"12500", "25000", "50000", "100000"})
+ private int totalTopicCount;
+ @Param({"10"})
+ private int partitionsPerTopic;
+ @Param({"3"})
+ private int replicationFactor;
+ @Param({"10000"})
+ private int numReplicasPerBroker;
+
+ private TopicsDelta topicsDelta;
+
+
+ @Setup(Level.Trial)
+ public void setup() {
+ // build a delta to apply within the benchmark code
+ // that consists of all the topics and partitions that would get loaded in a snapshot
+ topicsDelta = getInitialTopicsDelta(totalTopicCount, partitionsPerTopic, replicationFactor, numReplicasPerBroker);
+ System.out.print("(Loading a snapshot containing " + totalTopicCount + " total topics) ");
+ }
+
+ static TopicsDelta getInitialTopicsDelta(int totalTopicCount, int partitionsPerTopic, int replicationFactor, int numReplicasPerBroker) {
+ int numBrokers = getNumBrokers(totalTopicCount, partitionsPerTopic, replicationFactor, numReplicasPerBroker);
+ TopicsDelta buildupTopicsDelta = new TopicsDelta(TopicsImage.EMPTY);
+ final AtomicInteger currentLeader = new AtomicInteger(0);
+ IntStream.range(0, totalTopicCount).forEach(topicNumber -> {
+ Uuid topicId = Uuid.randomUuid();
+ buildupTopicsDelta.replay(new TopicRecord().setName("topic" + topicNumber).setTopicId(topicId));
+ IntStream.range(0, partitionsPerTopic).forEach(partitionNumber -> {
+ ArrayList replicas = getReplicas(totalTopicCount, partitionsPerTopic, replicationFactor, numReplicasPerBroker, currentLeader.get());
+ ArrayList isr = new ArrayList<>(replicas);
+ buildupTopicsDelta.replay(new PartitionRecord().
+ setPartitionId(partitionNumber).
+ setTopicId(topicId).
+ setReplicas(replicas).
+ setIsr(isr).
+ setRemovingReplicas(Collections.emptyList()).
+ setAddingReplicas(Collections.emptyList()).
+ setLeader(currentLeader.get()));
+ currentLeader.set((1 + currentLeader.get()) % numBrokers);
+ });
+ });
+ return buildupTopicsDelta;
+ }
+
+ static ArrayList getReplicas(int totalTopicCount, int partitionsPerTopic, int replicationFactor, int numReplicasPerBroker, int currentLeader) {
+ ArrayList replicas = new ArrayList<>();
+ int numBrokers = getNumBrokers(totalTopicCount, partitionsPerTopic, replicationFactor, numReplicasPerBroker);
+ IntStream.range(0, replicationFactor).forEach(replicaNumber ->
+ replicas.add((replicaNumber + currentLeader) % numBrokers));
+ return replicas;
+ }
+
+ static int getNumBrokers(int totalTopicCount, int partitionsPerTopic, int replicationFactor, int numReplicasPerBroker) {
+ int numBrokers = totalTopicCount * partitionsPerTopic * replicationFactor / numReplicasPerBroker;
+ return numBrokers - numBrokers % 3;
+ }
+
+ @Benchmark
+ public void testTopicsDeltaSnapshotLoad() {
+ topicsDelta.apply();
+ }
+}
diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/metadata/TopicsImageZonalOutageBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/metadata/TopicsImageZonalOutageBenchmark.java
new file mode 100644
index 0000000000000..5890763397ac6
--- /dev/null
+++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/metadata/TopicsImageZonalOutageBenchmark.java
@@ -0,0 +1,99 @@
+/*
+ * 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.jmh.metadata;
+
+import org.apache.kafka.common.Uuid;
+import org.apache.kafka.common.metadata.PartitionRecord;
+import org.apache.kafka.image.TopicsDelta;
+import org.apache.kafka.image.TopicsImage;
+import org.openjdk.jmh.annotations.Benchmark;
+import org.openjdk.jmh.annotations.BenchmarkMode;
+import org.openjdk.jmh.annotations.Fork;
+import org.openjdk.jmh.annotations.Level;
+import org.openjdk.jmh.annotations.Measurement;
+import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Param;
+import org.openjdk.jmh.annotations.Scope;
+import org.openjdk.jmh.annotations.Setup;
+import org.openjdk.jmh.annotations.State;
+import org.openjdk.jmh.annotations.Warmup;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+
+@State(Scope.Benchmark)
+@Fork(value = 1)
+@Warmup(iterations = 3)
+@Measurement(iterations = 5)
+@BenchmarkMode(Mode.AverageTime)
+@OutputTimeUnit(TimeUnit.MILLISECONDS)
+public class TopicsImageZonalOutageBenchmark {
+ @Param({"12500", "25000", "50000", "100000"})
+ private int totalTopicCount;
+ @Param({"10"})
+ private int partitionsPerTopic;
+ @Param({"3"})
+ private int replicationFactor;
+ @Param({"10000"})
+ private int numReplicasPerBroker;
+
+ private TopicsDelta topicsDelta;
+
+
+ @Setup(Level.Trial)
+ public void setup() {
+ // build an image containing all of the specified topics and partitions
+ TopicsDelta buildupTopicsDelta = TopicsImageSnapshotLoadBenchmark.getInitialTopicsDelta(totalTopicCount, partitionsPerTopic, replicationFactor, numReplicasPerBroker);
+ TopicsImage builtupTopicsImage = buildupTopicsDelta.apply();
+ // build a delta to apply within the benchmark code
+ // that perturbs all the topic-partitions for broker 0
+ // (as might happen in a zonal outage, one broker at a time, ultimately across 1/3 of the brokers in the cluster).
+ // It turns out that
+ topicsDelta = new TopicsDelta(builtupTopicsImage);
+ Set perturbedTopics = new HashSet<>();
+ builtupTopicsImage.topicsById().forEach((topicId, topicImage) ->
+ topicImage.partitions().forEach((partitionNumber, partitionRegistration) -> {
+ List newIsr = Arrays.stream(partitionRegistration.isr).boxed().filter(n -> n != 0).collect(Collectors.toList());
+ if (newIsr.size() < replicationFactor) {
+ perturbedTopics.add(topicId);
+ topicsDelta.replay(new PartitionRecord().
+ setPartitionId(partitionNumber).
+ setTopicId(topicId).
+ setReplicas(Arrays.stream(partitionRegistration.replicas).boxed().collect(Collectors.toList())).
+ setIsr(newIsr).
+ setRemovingReplicas(Collections.emptyList()).
+ setAddingReplicas(Collections.emptyList()).
+ setLeader(newIsr.get(0)));
+ }
+ })
+ );
+ int numBrokers = TopicsImageSnapshotLoadBenchmark.getNumBrokers(totalTopicCount, partitionsPerTopic, replicationFactor, numReplicasPerBroker);
+ System.out.print("(Perturbing 1 of " + numBrokers + " brokers, or " + perturbedTopics.size() + " topics within metadata having " + totalTopicCount + " total topics) ");
+ }
+
+ @Benchmark
+ public void testTopicsDeltaZonalOutage() {
+ topicsDelta.apply();
+ }
+}
diff --git a/licenses/pcollections-MIT b/licenses/pcollections-MIT
new file mode 100644
index 0000000000000..50519c5e432cd
--- /dev/null
+++ b/licenses/pcollections-MIT
@@ -0,0 +1,24 @@
+MIT License
+
+Copyright 2008-2011, 2014-2020, 2022 Harold Cooper, gil cattaneo, Gleb Frank,
+Günther Grill, Ilya Gorbunov, Jirka Kremser, Jochen Theodorou, Johnny Lim,
+Liam Miller, Mark Perry, Matei Dragu, Mike Klein, Oleg Osipenko, Ran Ari-Gur,
+Shantanu Kumar, and Valeriy Vyrva.
+
+Permission is hereby granted, free of charge, to any person obtaining a copy
+of this software and associated documentation files (the "Software"), to deal
+in the Software without restriction, including without limitation the rights
+to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
+copies of the Software, and to permit persons to whom the Software is
+furnished to do so, subject to the following conditions:
+
+The above copyright notice and this permission notice shall be included in
+all copies or substantial portions of the Software.
+
+THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
+IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
+FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
+AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
+LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
+OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
+SOFTWARE.
diff --git a/metadata/src/main/java/org/apache/kafka/image/TopicsDelta.java b/metadata/src/main/java/org/apache/kafka/image/TopicsDelta.java
index d3c5888fa7a2b..3927e19191dad 100644
--- a/metadata/src/main/java/org/apache/kafka/image/TopicsDelta.java
+++ b/metadata/src/main/java/org/apache/kafka/image/TopicsDelta.java
@@ -24,13 +24,13 @@
import org.apache.kafka.common.metadata.RemoveTopicRecord;
import org.apache.kafka.common.metadata.TopicRecord;
import org.apache.kafka.metadata.Replicas;
+import org.apache.kafka.server.immutable.ImmutableMap;
import org.apache.kafka.server.common.MetadataVersion;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Map;
-import java.util.Map.Entry;
import java.util.Set;
@@ -126,29 +126,27 @@ public void handleMetadataVersionChange(MetadataVersion newVersion) {
}
public TopicsImage apply() {
- Map newTopicsById = new HashMap<>(image.topicsById().size());
- Map newTopicsByName = new HashMap<>(image.topicsByName().size());
- for (Entry entry : image.topicsById().entrySet()) {
- Uuid id = entry.getKey();
- TopicImage prevTopicImage = entry.getValue();
- TopicDelta delta = changedTopics.get(id);
- if (delta == null) {
- if (!deletedTopicIds.contains(id)) {
- newTopicsById.put(id, prevTopicImage);
- newTopicsByName.put(prevTopicImage.name(), prevTopicImage);
- }
+ ImmutableMap newTopicsById = image.topicsById();
+ ImmutableMap newTopicsByName = image.topicsByName();
+ // apply all the deletes
+ for (Uuid topicId: deletedTopicIds) {
+ // it was deleted, so we have to remove it from the maps
+ TopicImage originalTopicToBeDeleted = image.topicsById().get(topicId);
+ if (originalTopicToBeDeleted == null) {
+ throw new IllegalStateException("Missing topic id " + topicId);
} else {
- TopicImage newTopicImage = delta.apply();
- newTopicsById.put(id, newTopicImage);
- newTopicsByName.put(delta.name(), newTopicImage);
+ newTopicsById = newTopicsById.removed(topicId);
+ newTopicsByName = newTopicsByName.removed(originalTopicToBeDeleted.name());
}
}
- for (Entry entry : changedTopics.entrySet()) {
- if (!newTopicsById.containsKey(entry.getKey())) {
- TopicImage newTopicImage = entry.getValue().apply();
- newTopicsById.put(newTopicImage.id(), newTopicImage);
- newTopicsByName.put(newTopicImage.name(), newTopicImage);
- }
+ // apply all the updates/additions
+ for (Map.Entry entry: changedTopics.entrySet()) {
+ Uuid topicId = entry.getKey();
+ TopicImage newTopicToBeAddedOrUpdated = entry.getValue().apply();
+ // put new information into the maps
+ String topicName = newTopicToBeAddedOrUpdated.name();
+ newTopicsById = newTopicsById.updated(topicId, newTopicToBeAddedOrUpdated);
+ newTopicsByName = newTopicsByName.updated(topicName, newTopicToBeAddedOrUpdated);
}
return new TopicsImage(newTopicsById, newTopicsByName);
}
diff --git a/metadata/src/main/java/org/apache/kafka/image/TopicsImage.java b/metadata/src/main/java/org/apache/kafka/image/TopicsImage.java
index 5f7db112f0aca..569264b1c4cb1 100644
--- a/metadata/src/main/java/org/apache/kafka/image/TopicsImage.java
+++ b/metadata/src/main/java/org/apache/kafka/image/TopicsImage.java
@@ -21,41 +21,45 @@
import org.apache.kafka.image.writer.ImageWriter;
import org.apache.kafka.image.writer.ImageWriterOptions;
import org.apache.kafka.metadata.PartitionRegistration;
+import org.apache.kafka.server.immutable.ImmutableMap;
import org.apache.kafka.server.util.TranslatedValueMapView;
-import java.util.Collections;
import java.util.Map;
import java.util.Objects;
import java.util.stream.Collectors;
-
/**
* Represents the topics in the metadata image.
*
* This class is thread-safe.
*/
public final class TopicsImage {
- public static final TopicsImage EMPTY =
- new TopicsImage(Collections.emptyMap(), Collections.emptyMap());
+ public static final TopicsImage EMPTY = new TopicsImage(ImmutableMap.empty(), ImmutableMap.empty());
+
+ private final ImmutableMap topicsById;
+ private final ImmutableMap topicsByName;
- private final Map topicsById;
- private final Map topicsByName;
+ public TopicsImage(ImmutableMap topicsById,
+ ImmutableMap topicsByName) {
+ this.topicsById = topicsById;
+ this.topicsByName = topicsByName;
+ }
- public TopicsImage(Map topicsById,
- Map topicsByName) {
- this.topicsById = Collections.unmodifiableMap(topicsById);
- this.topicsByName = Collections.unmodifiableMap(topicsByName);
+ public TopicsImage including(TopicImage topic) {
+ return new TopicsImage(
+ this.topicsById.updated(topic.id(), topic),
+ this.topicsByName.updated(topic.name(), topic));
}
public boolean isEmpty() {
return topicsById.isEmpty() && topicsByName.isEmpty();
}
- public Map topicsById() {
+ public ImmutableMap topicsById() {
return topicsById;
}
- public Map topicsByName() {
+ public ImmutableMap topicsByName() {
return topicsByName;
}
@@ -74,8 +78,8 @@ public TopicImage getTopic(String name) {
}
public void write(ImageWriter writer, ImageWriterOptions options) {
- for (TopicImage topicImage : topicsById.values()) {
- topicImage.write(writer, options);
+ for (Map.Entry entry : topicsById.entrySet()) {
+ entry.getValue().write(writer, options);
}
}
diff --git a/metadata/src/test/java/org/apache/kafka/controller/metrics/ControllerMetricsTestUtils.java b/metadata/src/test/java/org/apache/kafka/controller/metrics/ControllerMetricsTestUtils.java
index 7781bbdce9fc3..8be8548433a10 100644
--- a/metadata/src/test/java/org/apache/kafka/controller/metrics/ControllerMetricsTestUtils.java
+++ b/metadata/src/test/java/org/apache/kafka/controller/metrics/ControllerMetricsTestUtils.java
@@ -99,12 +99,10 @@ public static TopicImage fakeTopicImage(
public static TopicsImage fakeTopicsImage(
TopicImage... topics
) {
- Map topicsById = new HashMap<>();
- Map topicsByName = new HashMap<>();
+ TopicsImage image = TopicsImage.EMPTY;
for (TopicImage topic : topics) {
- topicsById.put(topic.id(), topic);
- topicsByName.put(topic.name(), topic);
+ image = image.including(topic);
}
- return new TopicsImage(topicsById, topicsByName);
+ return image;
}
}
diff --git a/metadata/src/test/java/org/apache/kafka/image/TopicsImageTest.java b/metadata/src/test/java/org/apache/kafka/image/TopicsImageTest.java
index 8268ec8609184..11af3488bf568 100644
--- a/metadata/src/test/java/org/apache/kafka/image/TopicsImageTest.java
+++ b/metadata/src/test/java/org/apache/kafka/image/TopicsImageTest.java
@@ -29,6 +29,7 @@
import org.apache.kafka.metadata.PartitionRegistration;
import org.apache.kafka.metadata.RecordTestUtils;
import org.apache.kafka.metadata.Replicas;
+import org.apache.kafka.server.immutable.ImmutableMap;
import org.apache.kafka.server.common.ApiMessageAndVersion;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
@@ -74,18 +75,18 @@ private static TopicImage newTopicImage(String name, Uuid id, PartitionRegistrat
return new TopicImage(name, id, partitionMap);
}
- private static Map newTopicsByIdMap(Collection topics) {
- Map map = new HashMap<>();
+ private static ImmutableMap newTopicsByIdMap(Collection topics) {
+ ImmutableMap map = TopicsImage.EMPTY.topicsById();
for (TopicImage topic : topics) {
- map.put(topic.id(), topic);
+ map = map.updated(topic.id(), topic);
}
return map;
}
- private static Map newTopicsByNameMap(Collection topics) {
- Map map = new HashMap<>();
+ private static ImmutableMap newTopicsByNameMap(Collection topics) {
+ ImmutableMap map = TopicsImage.EMPTY.topicsByName();
for (TopicImage topic : topics) {
- map.put(topic.name(), topic);
+ map = map.updated(topic.name(), topic);
}
return map;
}
diff --git a/server-common/src/main/java/org/apache/kafka/server/immutable/ImmutableMap.java b/server-common/src/main/java/org/apache/kafka/server/immutable/ImmutableMap.java
new file mode 100644
index 0000000000000..ec9790b3ff3ea
--- /dev/null
+++ b/server-common/src/main/java/org/apache/kafka/server/immutable/ImmutableMap.java
@@ -0,0 +1,64 @@
+/*
+ * 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.server.immutable;
+
+import org.apache.kafka.server.immutable.pcollections.PCollectionsImmutableMap;
+
+import java.util.Map;
+
+/**
+ * A persistent Hash-based Map wrapper.
+ * java.util.Map methods that mutate in-place will throw UnsupportedOperationException
+ *
+ * @param the key type
+ * @param the value type
+ */
+public interface ImmutableMap extends Map {
+ /**
+ * @return a wrapped hash-based persistent map that is empty
+ * @param the key type
+ * @param the value type
+ */
+ static ImmutableMap empty() {
+ return PCollectionsImmutableMap.empty();
+ }
+
+ /**
+ * @param key the key
+ * @param value the value
+ * @return a wrapped hash-based persistent map that has a single mapping
+ * @param the key type
+ * @param the value type
+ */
+ static ImmutableMap singleton(K key, V value) {
+ return PCollectionsImmutableMap.singleton(key, value);
+ }
+
+ /**
+ * @param key the key
+ * @param value the value
+ * @return a wrapped persistent map that differs from this one in that the given mapping is added (if necessary)
+ */
+ ImmutableMap updated(K key, V value);
+
+ /**
+ * @param key the key
+ * @return a wrapped persistent map that differs from this one in that the given mapping is removed (if necessary)
+ */
+ ImmutableMap removed(K key);
+}
diff --git a/server-common/src/main/java/org/apache/kafka/server/immutable/ImmutableSet.java b/server-common/src/main/java/org/apache/kafka/server/immutable/ImmutableSet.java
new file mode 100644
index 0000000000000..6f2cfd015a92f
--- /dev/null
+++ b/server-common/src/main/java/org/apache/kafka/server/immutable/ImmutableSet.java
@@ -0,0 +1,60 @@
+/*
+ * 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.server.immutable;
+
+import org.apache.kafka.server.immutable.pcollections.PCollectionsImmutableSet;
+
+import java.util.Set;
+
+/**
+ * A persistent Hash-based Set wrapper
+ * java.util.Set methods that mutate in-place will throw UnsupportedOperationException
+ *
+ * @param the element type
+ */
+public interface ImmutableSet extends Set {
+
+ /**
+ * @return a wrapped hash-based persistent set that is empty
+ * @param the element type
+ */
+ static ImmutableSet empty() {
+ return PCollectionsImmutableSet.empty();
+ }
+
+ /**
+ * @param e the element
+ * @return a wrapped hash-based persistent set that has a single element
+ * @param the element type
+ */
+ static ImmutableSet singleton(E e) {
+ return PCollectionsImmutableSet.singleton(e);
+ }
+
+ /**
+ * @param e the element
+ * @return a wrapped persistent set that differs from this one in that the given element is added (if necessary)
+ */
+ ImmutableSet added(E e);
+
+ /**
+ * @param e the element
+ * @return a wrapped persistent set that differs from this one in that the given element is added (if necessary)
+ */
+ ImmutableSet removed(E e);
+}
diff --git a/server-common/src/main/java/org/apache/kafka/server/immutable/pcollections/PCollectionsImmutableMap.java b/server-common/src/main/java/org/apache/kafka/server/immutable/pcollections/PCollectionsImmutableMap.java
new file mode 100644
index 0000000000000..d808f0db9c3fd
--- /dev/null
+++ b/server-common/src/main/java/org/apache/kafka/server/immutable/pcollections/PCollectionsImmutableMap.java
@@ -0,0 +1,223 @@
+/*
+ * 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.server.immutable.pcollections;
+
+import org.apache.kafka.server.immutable.ImmutableMap;
+import org.pcollections.HashPMap;
+import org.pcollections.HashTreePMap;
+
+import java.util.Collection;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+import java.util.function.BiConsumer;
+import java.util.function.BiFunction;
+import java.util.function.Function;
+
+@SuppressWarnings("deprecation")
+public class PCollectionsImmutableMap implements ImmutableMap {
+
+ private final HashPMap underlying;
+
+ /**
+ * @return a wrapped hash-based persistent map that is empty
+ * @param the key type
+ * @param the value type
+ */
+ public static PCollectionsImmutableMap empty() {
+ return new PCollectionsImmutableMap<>(HashTreePMap.empty());
+ }
+
+ /**
+ * @param key the key
+ * @param value the value
+ * @return a wrapped hash-based persistent map that has a single mapping
+ * @param the key type
+ * @param the value type
+ */
+ public static PCollectionsImmutableMap singleton(K key, V value) {
+ return new PCollectionsImmutableMap<>(HashTreePMap.singleton(key, value));
+ }
+
+ public PCollectionsImmutableMap(HashPMap map) {
+ this.underlying = Objects.requireNonNull(map);
+ }
+
+ @Override
+ public ImmutableMap updated(K key, V value) {
+ return new PCollectionsImmutableMap<>(underlying().plus(key, value));
+ }
+
+ @Override
+ public ImmutableMap removed(K key) {
+ return new PCollectionsImmutableMap<>(underlying().minus(key));
+ }
+
+ @Override
+ public int size() {
+ return underlying().size();
+ }
+
+ @Override
+ public boolean isEmpty() {
+ return underlying().isEmpty();
+ }
+
+ @Override
+ public boolean containsKey(Object key) {
+ return underlying().containsKey(key);
+ }
+
+ @Override
+ public boolean containsValue(Object value) {
+ return underlying().containsValue(value);
+ }
+
+ @Override
+ public V get(Object key) {
+ return underlying().get(key);
+ }
+
+ @Override
+ public V put(K key, V value) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().put(key, value);
+ }
+
+ @Override
+ public V remove(Object key) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().remove(key);
+ }
+
+ @Override
+ public void putAll(Map extends K, ? extends V> m) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ underlying().putAll(m);
+ }
+
+ @Override
+ public void clear() {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ underlying().clear();
+ }
+
+ @Override
+ public Set keySet() {
+ return underlying().keySet();
+ }
+
+ @Override
+ public Collection values() {
+ return underlying().values();
+ }
+
+ @Override
+ public Set> entrySet() {
+ return underlying().entrySet();
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) return true;
+ if (o == null || getClass() != o.getClass()) return false;
+ PCollectionsImmutableMap, ?> that = (PCollectionsImmutableMap, ?>) o;
+ return underlying().equals(that.underlying());
+ }
+
+ @Override
+ public int hashCode() {
+ return underlying().hashCode();
+ }
+
+ @Override
+ public V getOrDefault(Object key, V defaultValue) {
+ return underlying().getOrDefault(key, defaultValue);
+ }
+
+ @Override
+ public void forEach(BiConsumer super K, ? super V> action) {
+ underlying().forEach(action);
+ }
+
+ @Override
+ public void replaceAll(BiFunction super K, ? super V, ? extends V> function) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ underlying().replaceAll(function);
+ }
+
+ @Override
+ public V putIfAbsent(K key, V value) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().putIfAbsent(key, value);
+ }
+
+ @Override
+ public boolean remove(Object key, Object value) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().remove(key, value);
+ }
+
+ @Override
+ public boolean replace(K key, V oldValue, V newValue) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().replace(key, oldValue, newValue);
+ }
+
+ @Override
+ public V replace(K key, V value) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().replace(key, value);
+ }
+
+ @Override
+ public V computeIfAbsent(K key, Function super K, ? extends V> mappingFunction) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().computeIfAbsent(key, mappingFunction);
+ }
+
+ @Override
+ public V computeIfPresent(K key, BiFunction super K, ? super V, ? extends V> remappingFunction) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().computeIfPresent(key, remappingFunction);
+ }
+
+ @Override
+ public V compute(K key, BiFunction super K, ? super V, ? extends V> remappingFunction) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().compute(key, remappingFunction);
+ }
+
+ @Override
+ public V merge(K key, V value, BiFunction super V, ? super V, ? extends V> remappingFunction) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().merge(key, value, remappingFunction);
+ }
+
+ @Override
+ public String toString() {
+ return "PCollectionsImmutableMap{" +
+ "underlying=" + underlying() +
+ '}';
+ }
+
+ // package-private for testing
+ HashPMap underlying() {
+ return underlying;
+ }
+}
diff --git a/server-common/src/main/java/org/apache/kafka/server/immutable/pcollections/PCollectionsImmutableSet.java b/server-common/src/main/java/org/apache/kafka/server/immutable/pcollections/PCollectionsImmutableSet.java
new file mode 100644
index 0000000000000..8a50326ef1f7e
--- /dev/null
+++ b/server-common/src/main/java/org/apache/kafka/server/immutable/pcollections/PCollectionsImmutableSet.java
@@ -0,0 +1,188 @@
+/*
+ * 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.server.immutable.pcollections;
+
+import org.apache.kafka.server.immutable.ImmutableSet;
+import org.pcollections.HashTreePSet;
+import org.pcollections.MapPSet;
+
+import java.util.Collection;
+import java.util.Iterator;
+import java.util.Objects;
+import java.util.Spliterator;
+import java.util.function.Consumer;
+import java.util.function.Predicate;
+import java.util.stream.Stream;
+
+@SuppressWarnings("deprecation")
+public class PCollectionsImmutableSet implements ImmutableSet {
+ private final MapPSet underlying;
+
+ /**
+ * @return a wrapped hash-based persistent set that is empty
+ * @param the element type
+ */
+ public static PCollectionsImmutableSet empty() {
+ return new PCollectionsImmutableSet<>(HashTreePSet.empty());
+ }
+
+ /**
+ * @param e the element
+ * @return a wrapped hash-based persistent set that has a single element
+ * @param the element type
+ */
+ public static PCollectionsImmutableSet singleton(E e) {
+ return new PCollectionsImmutableSet<>(HashTreePSet.singleton(e));
+ }
+
+ public PCollectionsImmutableSet(MapPSet set) {
+ this.underlying = Objects.requireNonNull(set);
+ }
+
+ @Override
+ public ImmutableSet added(E e) {
+ return new PCollectionsImmutableSet<>(underlying().plus(e));
+ }
+
+ @Override
+ public ImmutableSet removed(E e) {
+ return new PCollectionsImmutableSet<>(underlying().minus(e));
+ }
+
+ @Override
+ public int size() {
+ return underlying().size();
+ }
+
+ @Override
+ public boolean isEmpty() {
+ return underlying().isEmpty();
+ }
+
+ @Override
+ public boolean contains(Object o) {
+ return underlying().contains(o);
+ }
+
+ @Override
+ public Iterator iterator() {
+ return underlying.iterator();
+ }
+
+ @Override
+ public void forEach(Consumer super E> action) {
+ underlying().forEach(action);
+ }
+
+ @Override
+ public Object[] toArray() {
+ return underlying().toArray();
+ }
+
+ @Override
+ public T[] toArray(T[] a) {
+ return underlying().toArray(a);
+ }
+
+ @Override
+ public boolean add(E e) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().add(e);
+ }
+
+ @Override
+ public boolean remove(Object o) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().remove(o);
+ }
+
+ @Override
+ public boolean containsAll(Collection> c) {
+ return underlying.containsAll(c);
+ }
+
+ @Override
+ public boolean addAll(Collection extends E> c) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().addAll(c);
+ }
+
+ @Override
+ public boolean retainAll(Collection> c) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().retainAll(c);
+ }
+
+ @Override
+ public boolean removeAll(Collection> c) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().removeAll(c);
+ }
+
+ @Override
+ public boolean removeIf(Predicate super E> filter) {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ return underlying().removeIf(filter);
+ }
+
+ @Override
+ public void clear() {
+ // will throw UnsupportedOperationException; delegate anyway for testability
+ underlying().clear();
+ }
+
+ @Override
+ public Spliterator spliterator() {
+ return underlying().spliterator();
+ }
+
+ @Override
+ public Stream stream() {
+ return underlying().stream();
+ }
+
+ @Override
+ public Stream parallelStream() {
+ return underlying().parallelStream();
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) return true;
+ if (o == null || getClass() != o.getClass()) return false;
+ PCollectionsImmutableSet> that = (PCollectionsImmutableSet>) o;
+ return Objects.equals(underlying(), that.underlying());
+ }
+
+ @Override
+ public int hashCode() {
+ return underlying().hashCode();
+ }
+
+ @Override
+ public String toString() {
+ return "PCollectionsImmutableSet{" +
+ "underlying=" + underlying() +
+ '}';
+ }
+
+ // package-private for testing
+ MapPSet underlying() {
+ return this.underlying;
+ }
+}
diff --git a/server-common/src/test/java/org/apache/kafka/server/immutable/DelegationChecker.java b/server-common/src/test/java/org/apache/kafka/server/immutable/DelegationChecker.java
new file mode 100644
index 0000000000000..e657e9b11c400
--- /dev/null
+++ b/server-common/src/test/java/org/apache/kafka/server/immutable/DelegationChecker.java
@@ -0,0 +1,146 @@
+/*
+ * 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.server.immutable;
+
+import org.mockito.Mockito;
+
+import java.util.Objects;
+import java.util.function.Consumer;
+import java.util.function.Function;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.when;
+
+/**
+ * Facilitate testing of wrapper class delegation.
+ *
+ * We require the following things to test delegation:
+ *
+ * 1. A mock object to which the wrapper is expected to delegate method invocations
+ * 2. A way to define how the mock is expected to behave when its method is invoked
+ * 3. A way to define how to invoke the method on the wrapper
+ * 4. A way to test that the method on the mock is invoked correctly when the wrapper method is invoked
+ * 5. A way to test that any return value from the wrapper method is correct
+
+ * @param delegate type
+ * @param wrapper type
+ * @param delegating method return type, if any
+ */
+public abstract class DelegationChecker {
+ private final D mock;
+ private final W wrapper;
+ private Consumer mockConsumer;
+ private Function mockConfigurationFunction;
+ private T mockFunctionReturnValue;
+ private Consumer wrapperConsumer;
+ private Function wrapperFunctionApplier;
+ private Function mockFunctionReturnValueTransformation;
+ private boolean expectWrapperToWrapMockFunctionReturnValue;
+ private boolean persistentCollectionMethodInvokedCorrectly = false;
+
+ /**
+ * @param mock mock for the underlying delegate
+ * @param wrapperCreator how to create a wrapper for the mock
+ */
+ protected DelegationChecker(D mock, Function wrapperCreator) {
+ this.mock = Objects.requireNonNull(mock);
+ this.wrapper = Objects.requireNonNull(wrapperCreator).apply(mock);
+ }
+
+ /**
+ * @param wrapper the wrapper
+ * @return the underlying delegate for the given wrapper
+ */
+ public abstract D unwrap(W wrapper);
+
+ public DelegationChecker defineMockConfigurationForVoidMethodInvocation(Consumer mockConsumer) {
+ this.mockConsumer = Objects.requireNonNull(mockConsumer);
+ return this;
+ }
+
+ public DelegationChecker defineMockConfigurationForFunctionInvocation(Function mockConfigurationFunction, T mockFunctionReturnValue) {
+ this.mockConfigurationFunction = Objects.requireNonNull(mockConfigurationFunction);
+ this.mockFunctionReturnValue = mockFunctionReturnValue;
+ return this;
+ }
+
+ public DelegationChecker defineWrapperVoidMethodInvocation(Consumer wrapperConsumer) {
+ this.wrapperConsumer = Objects.requireNonNull(wrapperConsumer);
+ return this;
+ }
+
+ public DelegationChecker defineWrapperFunctionInvocationAndMockReturnValueTransformation(
+ Function wrapperFunctionApplier,
+ Function expectedFunctionReturnValueTransformation) {
+ this.wrapperFunctionApplier = Objects.requireNonNull(wrapperFunctionApplier);
+ this.mockFunctionReturnValueTransformation = Objects.requireNonNull(expectedFunctionReturnValueTransformation);
+ return this;
+ }
+
+ public DelegationChecker expectWrapperToWrapMockFunctionReturnValue() {
+ this.expectWrapperToWrapMockFunctionReturnValue = true;
+ return this;
+ }
+
+ public void doVoidMethodDelegationCheck() {
+ if (mockConsumer == null || wrapperConsumer == null ||
+ mockConfigurationFunction != null || wrapperFunctionApplier != null ||
+ mockFunctionReturnValue != null || mockFunctionReturnValueTransformation != null) {
+ throwExceptionForIllegalTestSetup();
+ }
+ // configure the mock to behave as desired
+ mockConsumer.accept(Mockito.doAnswer(invocation -> {
+ persistentCollectionMethodInvokedCorrectly = true;
+ return null;
+ }).when(mock));
+ // invoke the wrapper, which should invoke the mock as desired
+ wrapperConsumer.accept(wrapper);
+ // assert that the expected delegation to the mock actually occurred
+ assertTrue(persistentCollectionMethodInvokedCorrectly);
+ }
+
+ @SuppressWarnings("unchecked")
+ public void doFunctionDelegationCheck() {
+ if (mockConfigurationFunction == null || wrapperFunctionApplier == null ||
+ mockFunctionReturnValueTransformation == null ||
+ mockConsumer != null || wrapperConsumer != null) {
+ throwExceptionForIllegalTestSetup();
+ }
+ // configure the mock to behave as desired
+ when(mockConfigurationFunction.apply(mock)).thenAnswer(invocation -> {
+ persistentCollectionMethodInvokedCorrectly = true;
+ return mockFunctionReturnValue;
+ });
+ // invoke the wrapper, which should invoke the mock as desired
+ T wrapperReturnValue = wrapperFunctionApplier.apply(wrapper);
+ // assert that the expected delegation to the mock actually occurred, including any return value transformation
+ assertTrue(persistentCollectionMethodInvokedCorrectly);
+ Object transformedMockFunctionReturnValue = mockFunctionReturnValueTransformation.apply(mockFunctionReturnValue);
+ if (this.expectWrapperToWrapMockFunctionReturnValue) {
+ assertEquals(transformedMockFunctionReturnValue, unwrap((W) wrapperReturnValue));
+ } else {
+ assertEquals(transformedMockFunctionReturnValue, wrapperReturnValue);
+ }
+ }
+
+ private static void throwExceptionForIllegalTestSetup() {
+ throw new IllegalStateException(
+ "test setup error: must define both mock and wrapper consumers or both mock and wrapper functions");
+ }
+}
diff --git a/server-common/src/test/java/org/apache/kafka/server/immutable/pcollections/PCollectionsImmutableMapTest.java b/server-common/src/test/java/org/apache/kafka/server/immutable/pcollections/PCollectionsImmutableMapTest.java
new file mode 100644
index 0000000000000..ab32b32be72d5
--- /dev/null
+++ b/server-common/src/test/java/org/apache/kafka/server/immutable/pcollections/PCollectionsImmutableMapTest.java
@@ -0,0 +1,310 @@
+/*
+ * 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.server.immutable.pcollections;
+
+import org.apache.kafka.server.immutable.DelegationChecker;
+import org.apache.kafka.server.immutable.ImmutableMap;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+import org.pcollections.HashPMap;
+import org.pcollections.HashTreePMap;
+
+import java.util.Collections;
+import java.util.function.BiConsumer;
+import java.util.function.BiFunction;
+import java.util.function.Function;
+
+import static java.util.function.Function.identity;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+
+@SuppressWarnings({"unchecked", "deprecation"})
+public class PCollectionsImmutableMapTest {
+ private static final HashPMap