From ac36a3ff49bb174dc5d7982098fab5dc4cd583c4 Mon Sep 17 00:00:00 2001 From: zifeihan Date: Sun, 1 Nov 2020 15:28:14 +0800 Subject: [PATCH] Add ThreadPoolExecutor for handle kafka message. (#5718) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * 1.Add ThreadPoolExecutor for handle kafka message. Co-authored-by: Daming Co-authored-by: 吴晟 Wu Sheng --- CHANGES.md | 1 + docs/en/setup/backend/backend-fetcher.md | 8 +-- .../setup/backend/configuration-vocabulary.md | 4 +- .../src/main/resources/application.yml | 6 +- .../kafka/KafkaFetcherHandlerRegister.java | 62 ++++++++++++++----- .../kafka/module/KafkaFetcherConfig.java | 5 ++ .../library/server/grpc/GRPCServer.java | 1 + .../{grpc => pool}/CustomThreadFactory.java | 4 +- test/e2e/e2e-service-provider/pom.xml | 10 +-- 9 files changed, 66 insertions(+), 35 deletions(-) rename oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/{grpc => pool}/CustomThreadFactory.java (94%) diff --git a/CHANGES.md b/CHANGES.md index 7fff1356c..1d9c585b6 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -16,6 +16,7 @@ Release Notes. * Support keeping collecting the slowly segments in the sampling mechanism. * Support choose files to active the meter analyzer. * Improve Kubernetes service registry for ALS analysis. +* Add the thread pool to the Kafka fetcher to increase the performance. #### UI diff --git a/docs/en/setup/backend/backend-fetcher.md b/docs/en/setup/backend/backend-fetcher.md index de5f1d9e1..0520d3d0e 100644 --- a/docs/en/setup/backend/backend-fetcher.md +++ b/docs/en/setup/backend/backend-fetcher.md @@ -8,12 +8,12 @@ prometheus-fetcher: selector: ${SW_PROMETHEUS_FETCHER:default} default: active: ${SW_PROMETHEUS_FETCHER_ACTIVE:false} -``` +``` ### Configuration file Prometheus fetcher is configured via a configuration file. The configuration file defines everything related to fetching services and their instances, as well as which rule files to load. - + OAP can load the configuration at bootstrap. If the new configuration is not well-formed, OAP fails to start up. The files are located at `$CLASSPATH/fetcher-prom-rules`. @@ -23,7 +23,7 @@ A full example can be found [here](../../../../oap-server/server-bootstrap/src/m Generic placeholders are defined as follows: - * ``: a duration This will parse a textual representation of a duration. The formats accepted are based on + * ``: a duration This will parse a textual representation of a duration. The formats accepted are based on the ISO-8601 duration format `PnDTnHnMn.nS` with days considered to be exactly 24 hours. * ``: a string matching the regular expression \[a-zA-Z_\]\[a-zA-Z0-9_\]* * ``: a string of unicode characters @@ -33,7 +33,7 @@ Generic placeholders are defined as follows: ```yaml # How frequently to fetch targets. -fetcherInterval: +fetcherInterval: # Per-fetch timeout when fetching this target. fetcherTimeout: # The HTTP resource path on which to fetch metrics from targets. diff --git a/docs/en/setup/backend/configuration-vocabulary.md b/docs/en/setup/backend/configuration-vocabulary.md index 9f3d9a43a..e5104e267 100644 --- a/docs/en/setup/backend/configuration-vocabulary.md +++ b/docs/en/setup/backend/configuration-vocabulary.md @@ -201,8 +201,10 @@ core|default|role|Option values, `Mixed/Receiver/Aggregator`. **Receiver** mode | - | - | isSharding | it was true when OAP Server in cluster. | SW_KAFKA_FETCHER_IS_SHARDING | false | | - | - | createTopicIfNotExist | If true, create the Kafka topic when it does not exist. | - | true | | - | - | partitions | The number of partitions for the topic being created. | SW_KAFKA_FETCHER_PARTITIONS | 3 | -| - | - | enableMeterSystem | To enable to fetch and handle [Meter System](backend-meter.md) data. | SW_KAFKA_FETCHER_ENABLE_METER_SYSTEM | false +| - | - | enableMeterSystem | To enable to fetch and handle [Meter System](backend-meter.md) data. | SW_KAFKA_FETCHER_ENABLE_METER_SYSTEM | false | | - | - | replicationFactor | The replication factor for each partition in the topic being created. | SW_KAFKA_FETCHER_PARTITIONS_FACTOR | 2 | +| - | - | kafkaHandlerThreadPoolSize | Pool size of kafka message handler executor. | SW_KAFKA_HANDLER_THREAD_POOL_SIZE | CPU core * 2 | +| - | - | kafkaHandlerThreadPoolQueueSize | The queue size of kafka message handler executor. | SW_KAFKA_HANDLER_THREAD_POOL_QUEUE_SIZE | 10000 | | - | - | topicNameOfMeters | Specifying Kafka topic name for Meter system data. | - | skywalking-meters | | - | - | topicNameOfMetrics | Specifying Kafka topic name for JVM Metrics data. | - | skywalking-metrics | | - | - | topicNameOfProfiling | Specifying Kafka topic name for Profiling data. | - | skywalking-profilings | diff --git a/oap-server/server-bootstrap/src/main/resources/application.yml b/oap-server/server-bootstrap/src/main/resources/application.yml index 721f5adf9..bdcf67ae6 100755 --- a/oap-server/server-bootstrap/src/main/resources/application.yml +++ b/oap-server/server-bootstrap/src/main/resources/application.yml @@ -71,8 +71,8 @@ core: gRPCPort: ${SW_CORE_GRPC_PORT:11800} maxConcurrentCallsPerConnection: ${SW_CORE_GRPC_MAX_CONCURRENT_CALL:0} maxMessageSize: ${SW_CORE_GRPC_MAX_MESSAGE_SIZE:0} - gRPCThreadPoolQueueSize: ${SW_CORE_GRPC_POOL_QUEUE_SIZE:0} - gRPCThreadPoolSize: ${SW_CORE_GRPC_THREAD_POOL_SIZE:0} + gRPCThreadPoolQueueSize: ${SW_CORE_GRPC_POOL_QUEUE_SIZE:-1} + gRPCThreadPoolSize: ${SW_CORE_GRPC_THREAD_POOL_SIZE:-1} gRPCSslEnabled: ${SW_CORE_GRPC_SSL_ENABLED:false} gRPCSslKeyPath: ${SW_CORE_GRPC_SSL_KEY_PATH:""} gRPCSslCertChainPath: ${SW_CORE_GRPC_SSL_CERT_CHAIN_PATH:""} @@ -269,6 +269,8 @@ kafka-fetcher: enableMeterSystem: ${SW_KAFKA_FETCHER_ENABLE_METER_SYSTEM:false} isSharding: ${SW_KAFKA_FETCHER_IS_SHARDING:false} consumePartitions: ${SW_KAFKA_FETCHER_CONSUME_PARTITIONS:""} + kafkaHandlerThreadPoolSize: ${SW_KAFKA_HANDLER_THREAD_POOL_SIZE:-1} + kafkaHandlerThreadPoolQueueSize: ${SW_KAFKA_HANDLER_THREAD_POOL_QUEUE_SIZE:-1} receiver-meter: selector: ${SW_RECEIVER_METER:default} diff --git a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/KafkaFetcherHandlerRegister.java b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/KafkaFetcherHandlerRegister.java index 16b7ee10c..3b324ed4f 100644 --- a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/KafkaFetcherHandlerRegister.java +++ b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/KafkaFetcherHandlerRegister.java @@ -20,15 +20,16 @@ package org.apache.skywalking.oap.server.analyzer.agent.kafka; import com.google.common.collect.ImmutableMap; import com.google.common.collect.Lists; -import io.netty.util.concurrent.DefaultThreadFactory; import java.time.Duration; import java.util.Iterator; import java.util.List; import java.util.Objects; import java.util.Properties; import java.util.Set; +import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.ExecutionException; -import java.util.concurrent.Executors; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.admin.AdminClient; @@ -42,9 +43,10 @@ import org.apache.kafka.common.serialization.BytesDeserializer; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.utils.Bytes; import org.apache.skywalking.apm.util.StringUtil; +import org.apache.skywalking.oap.server.analyzer.agent.kafka.module.KafkaFetcherConfig; import org.apache.skywalking.oap.server.analyzer.agent.kafka.provider.handler.KafkaHandler; import org.apache.skywalking.oap.server.library.module.ModuleStartException; -import org.apache.skywalking.oap.server.analyzer.agent.kafka.module.KafkaFetcherConfig; +import org.apache.skywalking.oap.server.library.server.pool.CustomThreadFactory; /** * Configuring and initializing a KafkaConsumer client as a dispatcher to delivery Kafka Message to registered handler by topic. @@ -60,8 +62,14 @@ public class KafkaFetcherHandlerRegister implements Runnable { private final KafkaFetcherConfig config; private final boolean isSharding; + private int threadPoolSize = Runtime.getRuntime().availableProcessors() * 2; + private int threadPoolQueueSize = 10000; + private final ThreadPoolExecutor executor; + private final boolean enableKafkaMessageAutoCommit; + public KafkaFetcherHandlerRegister(KafkaFetcherConfig config) throws ModuleStartException { this.config = config; + Properties properties = new Properties(); properties.putAll(config.getKafkaConsumerConfig()); properties.setProperty(ConsumerConfig.GROUP_ID_CONFIG, config.getGroupId()); @@ -92,11 +100,11 @@ public class KafkaFetcherHandlerRegister implements Runnable { if (!missedTopics.isEmpty()) { log.info("Topics" + missedTopics.toString() + " not exist."); List newTopicList = missedTopics.stream() - .map(topic -> new NewTopic( - topic, - config.getPartitions(), - (short) config.getReplicationFactor() - )).collect(Collectors.toList()); + .map(topic -> new NewTopic( + topic, + config.getPartitions(), + (short) config.getReplicationFactor() + )).collect(Collectors.toList()); try { adminClient.createTopics(newTopicList).all().get(); @@ -110,7 +118,22 @@ public class KafkaFetcherHandlerRegister implements Runnable { } else { isSharding = false; } + if (config.getKafkaHandlerThreadPoolSize() > 0) { + threadPoolSize = config.getKafkaHandlerThreadPoolSize(); + } + if (config.getKafkaHandlerThreadPoolQueueSize() > 0) { + threadPoolQueueSize = config.getKafkaHandlerThreadPoolQueueSize(); + } + + enableKafkaMessageAutoCommit = (boolean) properties.getOrDefault( + ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true); consumer = new KafkaConsumer<>(properties, new StringDeserializer(), new BytesDeserializer()); + executor = new ThreadPoolExecutor(threadPoolSize, threadPoolSize, + 60, TimeUnit.SECONDS, + new ArrayBlockingQueue(threadPoolQueueSize), + new CustomThreadFactory("KafkaConsumer"), + new ThreadPoolExecutor.CallerRunsPolicy() + ); } public void register(KafkaHandler handler) { @@ -126,22 +149,27 @@ public class KafkaFetcherHandlerRegister implements Runnable { consumer.subscribe(handlerMap.keySet()); } consumer.seekToEnd(consumer.assignment()); - Executors.newSingleThreadExecutor(new DefaultThreadFactory("KafkaConsumer")).submit(this); + executor.submit(this); } @Override public void run() { while (true) { - ConsumerRecords consumerRecords = consumer.poll(Duration.ofMillis(500L)); - if (!consumerRecords.isEmpty()) { - Iterator> iterator = consumerRecords.iterator(); - while (iterator.hasNext()) { - ConsumerRecord record = iterator.next(); - handlerMap.get(record.topic()).handle(record); + try { + ConsumerRecords consumerRecords = consumer.poll(Duration.ofMillis(500L)); + if (!consumerRecords.isEmpty()) { + Iterator> iterator = consumerRecords.iterator(); + while (iterator.hasNext()) { + ConsumerRecord record = iterator.next(); + executor.submit(() -> handlerMap.get(record.topic()).handle(record)); + } + if (!enableKafkaMessageAutoCommit) { + consumer.commitAsync(); + } } - consumer.commitAsync(); + } catch (Exception e) { + log.error("Kafka handle message error.", e); } } } - } diff --git a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/module/KafkaFetcherConfig.java b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/module/KafkaFetcherConfig.java index c7b6239ee..27e6a2571 100644 --- a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/module/KafkaFetcherConfig.java +++ b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/module/KafkaFetcherConfig.java @@ -80,4 +80,9 @@ public class KafkaFetcherConfig extends ModuleConfig { private String topicNameOfManagements = "skywalking-managements"; private String topicNameOfMeters = "skywalking-meters"; + + private int kafkaHandlerThreadPoolSize; + + private int kafkaHandlerThreadPoolQueueSize; + } diff --git a/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/GRPCServer.java b/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/GRPCServer.java index e29c3e40e..1e0d3c408 100644 --- a/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/GRPCServer.java +++ b/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/GRPCServer.java @@ -36,6 +36,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.skywalking.oap.server.library.server.Server; import org.apache.skywalking.oap.server.library.server.ServerException; import org.apache.skywalking.oap.server.library.server.grpc.ssl.DynamicSslContext; +import org.apache.skywalking.oap.server.library.server.pool.CustomThreadFactory; @Slf4j public class GRPCServer implements Server { diff --git a/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/CustomThreadFactory.java b/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/pool/CustomThreadFactory.java similarity index 94% rename from oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/CustomThreadFactory.java rename to oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/pool/CustomThreadFactory.java index 4cfe2e171..643ab4db7 100644 --- a/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/CustomThreadFactory.java +++ b/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/pool/CustomThreadFactory.java @@ -16,7 +16,7 @@ * */ -package org.apache.skywalking.oap.server.library.server.grpc; +package org.apache.skywalking.oap.server.library.server.pool; import java.util.concurrent.ThreadFactory; import java.util.concurrent.atomic.AtomicInteger; @@ -27,7 +27,7 @@ public class CustomThreadFactory implements ThreadFactory { private final AtomicInteger threadNumber = new AtomicInteger(1); private final String namePrefix; - CustomThreadFactory(String name) { + public CustomThreadFactory(String name) { SecurityManager s = System.getSecurityManager(); group = (s != null) ? s.getThreadGroup() : Thread.currentThread().getThreadGroup(); namePrefix = name + "-" + poolNumber.getAndIncrement() + "-thread-"; diff --git a/test/e2e/e2e-service-provider/pom.xml b/test/e2e/e2e-service-provider/pom.xml index 693eab25d..19c31c460 100644 --- a/test/e2e/e2e-service-provider/pom.xml +++ b/test/e2e/e2e-service-provider/pom.xml @@ -51,7 +51,7 @@ org.apache.skywalking apm-toolkit-micrometer-registry - 8.2.0-SNAPSHOT + 8.2.0 @@ -76,12 +76,4 @@ - - - - apache-snapshot - Apache Snapshot - https://repository.apache.org/content/repositories/snapshots - -