diff --git a/conf/broker.conf b/conf/broker.conf index 5a2611cecb8fe..95f0049ed7284 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -143,7 +143,9 @@ webServiceLogDetailedAddresses= # Number of threads to config Netty Acceptor. Default is 1 numAcceptorThreads= -# Number of threads to use for Netty IO. Default is set to 2 * Runtime.getRuntime().availableProcessors() +# Number of threads to use for Netty IO. Default is set to half the available processors, but no fewer than 8 unless +# that exceeds twice the available processors: Math.min(2 * n, Math.max(8, n / 2)) for n available processors. +# Increase it when the busiest pulsar-io thread stays busy and consumers' backlogs grow numIOThreads= # Number of threads to use for ordered executor. The ordered executor is used to operate with zookeeper, diff --git a/conf/standalone.conf b/conf/standalone.conf index b9e8facabe8d9..9ee816d1c4d9a 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -100,7 +100,9 @@ webServiceTrustXForwardedFor=false # Defaults to true when either webServiceHaProxyProtocolEnabled or webServiceTrustXForwardedFor is enabled. webServiceLogDetailedAddresses= -# Number of threads to use for Netty IO. Default is set to 2 * Runtime.getRuntime().availableProcessors() +# Number of threads to use for Netty IO. Default is set to half the available processors, but no fewer than 8 unless +# that exceeds twice the available processors: Math.min(2 * n, Math.max(8, n / 2)) for n available processors. +# Increase it when the busiest pulsar-io thread stays busy and consumers' backlogs grow numIOThreads= # Number of threads to use for ordered executor. The ordered executor is used to operate with zookeeper, diff --git a/deployment/terraform-ansible/templates/broker.conf b/deployment/terraform-ansible/templates/broker.conf index 8f329d0094229..5e2cc0fd49e93 100644 --- a/deployment/terraform-ansible/templates/broker.conf +++ b/deployment/terraform-ansible/templates/broker.conf @@ -56,7 +56,9 @@ advertisedAddress={{ hostvars[inventory_hostname].private_ip }} # The Default value is absent, the broker uses the first listener as the internal listener. # internalListenerName= -# Number of threads to use for Netty IO. Default is set to 2 * Runtime.getRuntime().availableProcessors() +# Number of threads to use for Netty IO. Default is set to half the available processors, but no fewer than 8 unless +# that exceeds twice the available processors: Math.min(2 * n, Math.max(8, n / 2)) for n available processors. +# Increase it when the busiest pulsar-io thread stays busy and consumers' backlogs grow numIOThreads= # Number of threads to use for HTTP requests processing. Default is set to 2 * Runtime.getRuntime().availableProcessors() diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index fad7fcd421c61..df0ec5e0b7f6f 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -337,9 +337,20 @@ public class ServiceConfiguration implements PulsarConfiguration { @FieldContext( category = CATEGORY_SERVER, doc = "Number of threads to use for Netty IO." - + " Default is set to `2 * Runtime.getRuntime().availableProcessors()`" + + " Default is set to half the available processors, but no fewer than 8 unless that exceeds twice the" + + " available processors: `Math.min(2 * n, Math.max(8, n / 2))` for n available processors." + + " Increase it when the busiest pulsar-io thread stays busy and consumers' backlogs grow" ) - private int numIOThreads = 2 * Runtime.getRuntime().availableProcessors(); + private int numIOThreads = defaultNumIOThreads(Runtime.getRuntime().availableProcessors()); + + /** + * The default number of Netty IO threads for the given number of available processors: half of them, but no fewer + * than 8 unless that exceeds twice the processors. Fewer, busier event loops batch more work per wakeup and + * system call, and the threads that hand work over to them, such as a managed ledger's, wake them less often. + */ + static int defaultNumIOThreads(int availableProcessors) { + return Math.min(2 * availableProcessors, Math.max(8, availableProcessors / 2)); + } @FieldContext( category = CATEGORY_SERVER, diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/ServiceConfigurationDefaultsTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/ServiceConfigurationDefaultsTest.java new file mode 100644 index 0000000000000..b307c08721078 --- /dev/null +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/ServiceConfigurationDefaultsTest.java @@ -0,0 +1,49 @@ +/* + * 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.pulsar.broker; + +import static org.testng.Assert.assertEquals; +import org.testng.annotations.Test; + +public class ServiceConfigurationDefaultsTest { + @Test + public void defaultNumIOThreads() { + // up to 4 processors, the previous default of twice the processors + assertEquals(ServiceConfiguration.defaultNumIOThreads(1), 2); + assertEquals(ServiceConfiguration.defaultNumIOThreads(2), 4); + assertEquals(ServiceConfiguration.defaultNumIOThreads(3), 6); + assertEquals(ServiceConfiguration.defaultNumIOThreads(4), 8); + // 8 from 4 up to 17 processors + assertEquals(ServiceConfiguration.defaultNumIOThreads(5), 8); + assertEquals(ServiceConfiguration.defaultNumIOThreads(7), 8); + assertEquals(ServiceConfiguration.defaultNumIOThreads(8), 8); + assertEquals(ServiceConfiguration.defaultNumIOThreads(16), 8); + assertEquals(ServiceConfiguration.defaultNumIOThreads(17), 8); + // half the processors from 18 up + assertEquals(ServiceConfiguration.defaultNumIOThreads(18), 9); + assertEquals(ServiceConfiguration.defaultNumIOThreads(32), 16); + assertEquals(ServiceConfiguration.defaultNumIOThreads(128), 64); + } + + @Test + public void numIOThreadsDefaultsToDefaultNumIOThreadsOfTheAvailableProcessors() { + assertEquals(new ServiceConfiguration().getNumIOThreads(), + ServiceConfiguration.defaultNumIOThreads(Runtime.getRuntime().availableProcessors())); + } +}