Skip to content
Merged
Show file tree
Hide file tree
Changes from 13 commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
7afe1e3
KAFKA-14735: Improve KRaft metadata image change performance at high …
rondagostino Feb 5, 2023
f195c20
Use vavr instead of Paguro
rondagostino Feb 23, 2023
d08b81b
Use PCollections instead of vavr
rondagostino Mar 25, 2023
6ef59b3
Add wrapper classes
rondagostino Apr 3, 2023
d6c76c6
Merge remote-tracking branch 'apache/trunk' into KAFKA-14735
rondagostino Apr 3, 2023
a6887e2
cleanups
rondagostino Apr 3, 2023
634a9e7
more cleanups
rondagostino Apr 3, 2023
9140e66
Fix license check failure
rondagostino Apr 4, 2023
009dd63
Merge remote-tracking branch 'apache/trunk' into KAFKA-14735
rondagostino Apr 4, 2023
25947d8
Test PCollectionsHashSetWrapper.hashCode(), equals(), and toString()
rondagostino Apr 4, 2023
849fd66
Test PCollectionsHashSetWrapper.forEach()
rondagostino Apr 4, 2023
baa60d6
Better clarity on testing
rondagostino Apr 5, 2023
a61408e
cleanup
rondagostino Apr 5, 2023
7839943
Wrapper classes implement standard java interfaces
rondagostino Apr 7, 2023
c4667a6
Rename package pcoll to server.immutable
rondagostino Apr 8, 2023
19ad442
More fine-grained import control
rondagostino Apr 8, 2023
c2e1467
Rename classes
rondagostino Apr 8, 2023
b711e2a
Rename methods
rondagostino Apr 8, 2023
8b2a177
Merge remote-tracking branch 'apache/trunk' into KAFKA-14735
rondagostino Apr 8, 2023
6b24920
Make underlying() non-public
rondagostino Apr 10, 2023
79078d7
Eliminate factory in favor of static methods on interfaces
rondagostino Apr 11, 2023
323e771
Merge remote-tracking branch 'apache/trunk' into KAFKA-14735
rondagostino Apr 11, 2023
2515cab
Eliminate direct field access
rondagostino Apr 11, 2023
3d4f03e
Merge remote-tracking branch 'apache/trunk' into KAFKA-14735
rondagostino Apr 12, 2023
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion LICENSE-binary
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -1551,11 +1551,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
Expand Down
1 change: 1 addition & 0 deletions checkstyle/import-control-jmh-benchmarks.xml
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@
<allow pkg="org.apache.kafka.storage"/>
<allow pkg="org.apache.kafka.clients"/>
<allow pkg="org.apache.kafka.coordinator.group"/>
<allow pkg="org.apache.kafka.image"/>
<allow pkg="org.apache.kafka.metadata"/>
<allow pkg="org.apache.kafka.timeline" />

Expand Down
12 changes: 12 additions & 0 deletions checkstyle/import-control.xml
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,8 @@
<allow pkg="org.apache.kafka.common.utils" />
<allow pkg="org.apache.kafka.common.errors" exact-match="true" />
<allow pkg="org.apache.kafka.common.memory" />
<!-- anyone can use persistent collection factories/non-library-specific wrappers -->
<allow pkg="org.apache.kafka.pcoll" exact-match="true" />

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Well, "anyone" on the server side :) I don't think we want the kafka client to take this dependency. At least not yet, until we have a very clear use-case in mind.

So with that in mind, sadly it might be better to add this "allow" to server-common, metadata, and core individually, rather than up here. (although that's slightly more work I realize)


<subpackage name="common">
<allow class="org.apache.kafka.clients.consumer.ConsumerRecord" exact-match="true" />
Expand Down Expand Up @@ -317,6 +319,16 @@
<allow pkg="org.apache.kafka.test" />
</subpackage>

<subpackage name="pcoll">
<allow pkg="org.apache.kafka.server.util"/>
<!-- only the factory package can use persistent collection library-specific wrapper implementations -->
<!-- the library-specific wrapper implementation for PCollections -->
<allow pkg="org.apache.kafka.pcoll.pcollections" />
<subpackage name="pcollections">
<allow pkg="org.pcollections" />
</subpackage>
</subpackage>

<subpackage name="queue">
<allow pkg="org.apache.kafka.test" />
</subpackage>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
2 changes: 2 additions & 0 deletions gradle/dependencies.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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",
Expand Down
Original file line number Diff line number Diff line change
@@ -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<UpdateMetadataEndpoint> 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();
}
}
Loading