From ec2ce5603d8b5566781f6a0a720aca0d5c660789 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Tue, 10 Dec 2024 16:06:11 -0500 Subject: [PATCH 1/9] proof of concept --- .../kafka/controller/QuorumController.java | 6 ++ .../metrics/QuorumControllerMetrics.java | 10 ++++ .../controller/metrics/SlowEventsLogger.java | 58 +++++++++++++++++++ 3 files changed, 74 insertions(+) create mode 100644 metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java diff --git a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java index aaffb9084efbc..eb3ff7bf7007f 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -94,6 +94,7 @@ import org.apache.kafka.controller.errors.ControllerExceptions; import org.apache.kafka.controller.errors.EventHandlerExceptionInfo; import org.apache.kafka.controller.metrics.QuorumControllerMetrics; +import org.apache.kafka.controller.metrics.SlowEventsLogger; import org.apache.kafka.deferred.DeferredEvent; import org.apache.kafka.deferred.DeferredEventQueue; import org.apache.kafka.metadata.BrokerHeartbeatReply; @@ -151,6 +152,7 @@ import static java.util.concurrent.TimeUnit.MICROSECONDS; import static java.util.concurrent.TimeUnit.MILLISECONDS; import static java.util.concurrent.TimeUnit.NANOSECONDS; +import static java.util.concurrent.TimeUnit.SECONDS; import static org.apache.kafka.controller.QuorumController.ControllerOperationFlag.DOES_NOT_UPDATE_QUEUE_TIME; @@ -525,6 +527,7 @@ private void handleEventEnd(String name, long startProcessingTimeNs) { long deltaNs = endProcessingTime - startProcessingTimeNs; log.debug("Processed {} in {} us", name, MICROSECONDS.convert(deltaNs, NANOSECONDS)); + slowEventsLogger.maybeLog(name, deltaNs); controllerMetrics.updateEventQueueProcessingTime(NANOSECONDS.toMillis(deltaNs)); } @@ -1446,6 +1449,8 @@ private void replay(ApiMessage message, Optional snapshotId, lon */ private final RecordRedactor recordRedactor; + private final SlowEventsLogger slowEventsLogger; + private QuorumController( FaultHandler nonFatalFaultHandler, FaultHandler fatalFaultHandler, @@ -1599,6 +1604,7 @@ private QuorumController( log.info("Creating new QuorumController with clusterId {}", clusterId); this.raftClient.register(metaLogListener); + this.slowEventsLogger = new SlowEventsLogger(controllerMetrics::getEventQueueProcessingTimeP99, log, time); } /** diff --git a/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java b/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java index 10c3807da854b..3d43bf1440c61 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java +++ b/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java @@ -29,6 +29,8 @@ import java.util.Optional; import java.util.concurrent.atomic.AtomicLong; import java.util.function.Consumer; +import java.util.function.Function; +import java.util.function.Supplier; /** * These are the metrics which are managed by the QuorumController class. They generally pertain to @@ -164,6 +166,14 @@ public void updateEventQueueProcessingTime(long durationMs) { eventQueueProcessingTimeUpdater.accept(durationMs); } + public double getEventQueueProcessingTimeP99() { + if (registry.isPresent()) { + Histogram histogram = registry.get().newHistogram(EVENT_QUEUE_PROCESSING_TIME_MS, false); + return histogram.getSnapshot().get99thPercentile(); + } else { + return -1.0; + } + } public void setLastAppliedRecordOffset(long offset) { lastAppliedRecordOffset.set(offset); } diff --git a/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java b/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java new file mode 100644 index 0000000000000..eb694566062f2 --- /dev/null +++ b/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java @@ -0,0 +1,58 @@ +package org.apache.kafka.controller.metrics; + +import org.apache.kafka.common.utils.Time; +import org.apache.kafka.common.utils.Timer; +import org.slf4j.Logger; + +import java.util.function.Supplier; + +import static java.util.concurrent.TimeUnit.MICROSECONDS; +import static java.util.concurrent.TimeUnit.MILLISECONDS; +import static java.util.concurrent.TimeUnit.NANOSECONDS; + +public class SlowEventsLogger { + /** + * Don't report any p99 events below this threshold. This prevents the controller from reporting p99 event + * times in the idle case. + */ + private static final int MIN_SLOW_EVENT_TIME_MS = 100; + + /** + * Calculating the p99 from the histogram consumes some resources, so only update it every so often. + */ + private static final int P99_REFRESH_INTERVAL_MS = 30000; + + private final Supplier p99Supplier; + private final Logger log; + private final Time time; + private double p99; + private Timer percentileUpdateTimer; + + public SlowEventsLogger( + Supplier p99Supplier, + Logger log, + Time time + ) { + this.p99Supplier = p99Supplier; + this.log = log; + this.time = time; + this.percentileUpdateTimer = time.timer(P99_REFRESH_INTERVAL_MS); + } + + public void maybeLog(String name, long durationNs) { + long durationMs = MILLISECONDS.convert(durationNs, NANOSECONDS); + if (durationMs > MIN_SLOW_EVENT_TIME_MS && durationMs > p99) { + log.info("Slow controller event {} processed in {} us", name, + MICROSECONDS.convert(durationNs, NANOSECONDS)); + } + } + + public void maybeRefreshPercentile() { + percentileUpdateTimer.update(); + if (percentileUpdateTimer.isExpired()) { + p99 = p99Supplier.get(); + log.trace("Update slow events p99 to {}", p99); + percentileUpdateTimer.reset(P99_REFRESH_INTERVAL_MS); + } + } +} From 238ce5af63372c59790dd8444440c0c99ced0163 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Thu, 28 Mar 2024 19:39:50 -0400 Subject: [PATCH 2/9] wip --- .../kafka/controller/metrics/QuorumControllerMetrics.java | 3 ++- .../org/apache/kafka/controller/metrics/SlowEventsLogger.java | 3 +-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java b/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java index 3d43bf1440c61..ada68a4ff210c 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java +++ b/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java @@ -171,7 +171,8 @@ public double getEventQueueProcessingTimeP99() { Histogram histogram = registry.get().newHistogram(EVENT_QUEUE_PROCESSING_TIME_MS, false); return histogram.getSnapshot().get99thPercentile(); } else { - return -1.0; + // Only returned in unit tests when a metrics registry is not set. + return 0.0; } } public void setLastAppliedRecordOffset(long offset) { diff --git a/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java b/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java index eb694566062f2..857eec6c8d8ac 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java +++ b/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java @@ -24,7 +24,6 @@ public class SlowEventsLogger { private final Supplier p99Supplier; private final Logger log; - private final Time time; private double p99; private Timer percentileUpdateTimer; @@ -35,8 +34,8 @@ public SlowEventsLogger( ) { this.p99Supplier = p99Supplier; this.log = log; - this.time = time; this.percentileUpdateTimer = time.timer(P99_REFRESH_INTERVAL_MS); + this.p99 = p99Supplier.get(); } public void maybeLog(String name, long durationNs) { From af0482857b9f3bac77e44612ddf9544a3a4ae2a1 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Tue, 10 Dec 2024 21:24:17 -0500 Subject: [PATCH 3/9] WIP --- .../kafka/controller/QuorumController.java | 16 +++++++-- .../controller/metrics/SlowEventsLogger.java | 36 ++++++++----------- .../metrics/SlowEventsLoggerTest.java | 17 +++++++++ 3 files changed, 44 insertions(+), 25 deletions(-) create mode 100644 metadata/src/test/java/org/apache/kafka/controller/metrics/SlowEventsLoggerTest.java diff --git a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java index eb3ff7bf7007f..9c3c934a75a50 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -527,7 +527,7 @@ private void handleEventEnd(String name, long startProcessingTimeNs) { long deltaNs = endProcessingTime - startProcessingTimeNs; log.debug("Processed {} in {} us", name, MICROSECONDS.convert(deltaNs, NANOSECONDS)); - slowEventsLogger.maybeLog(name, deltaNs); + slowEventsLogger.maybeLogEvent(name, deltaNs); controllerMetrics.updateEventQueueProcessingTime(NANOSECONDS.toMillis(deltaNs)); } @@ -1592,7 +1592,7 @@ private QuorumController( } registerElectUnclean(TimeUnit.MILLISECONDS.toNanos(uncleanLeaderElectionCheckIntervalMs)); registerExpireDelegationTokens(MILLISECONDS.toNanos(delegationTokenExpiryCheckIntervalMs)); - + registerSlowEventUpdater(MILLISECONDS.toNanos(30000)); // OffsetControlManager must be initialized last, because its constructor will take the // initial in-memory snapshot of all extant timeline data structures. this.offsetControl = new OffsetControlManager.Builder(). @@ -1604,7 +1604,7 @@ private QuorumController( log.info("Creating new QuorumController with clusterId {}", clusterId); this.raftClient.register(metaLogListener); - this.slowEventsLogger = new SlowEventsLogger(controllerMetrics::getEventQueueProcessingTimeP99, log, time); + this.slowEventsLogger = new SlowEventsLogger(controllerMetrics::getEventQueueProcessingTimeP99, logContext); } /** @@ -1628,6 +1628,16 @@ private void registerWriteNoOpRecord(long maxIdleIntervalNs) { EnumSet.noneOf(PeriodicTaskFlag.class))); } + private void registerSlowEventUpdater(long maxSlowEventWindowNs) { + periodicControl.registerTask(new PeriodicTask("updateSlowEventP99", + () -> { + slowEventsLogger.refreshPercentile(); + return ControllerResult.of(Collections.emptyList(), false); + }, + maxSlowEventWindowNs, + EnumSet.noneOf(PeriodicTaskFlag.class))); + } + /** * Calculate what the period should be for the maybeFenceStaleBroker task. * diff --git a/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java b/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java index 857eec6c8d8ac..fb2a128e50d3e 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java +++ b/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java @@ -1,7 +1,7 @@ package org.apache.kafka.controller.metrics; -import org.apache.kafka.common.utils.Time; -import org.apache.kafka.common.utils.Timer; +import org.apache.kafka.common.utils.LogContext; +import org.apache.kafka.controller.QuorumController; import org.slf4j.Logger; import java.util.function.Supplier; @@ -10,6 +10,10 @@ import static java.util.concurrent.TimeUnit.MILLISECONDS; import static java.util.concurrent.TimeUnit.NANOSECONDS; +/** + * Track the p99 for controller event queue processing time. If we encounter an event that takes longer + * than this cached p99 time, we will log it at INFO level on the controller logger. + */ public class SlowEventsLogger { /** * Don't report any p99 events below this threshold. This prevents the controller from reporting p99 event @@ -17,41 +21,29 @@ public class SlowEventsLogger { */ private static final int MIN_SLOW_EVENT_TIME_MS = 100; - /** - * Calculating the p99 from the histogram consumes some resources, so only update it every so often. - */ - private static final int P99_REFRESH_INTERVAL_MS = 30000; - private final Supplier p99Supplier; private final Logger log; private double p99; - private Timer percentileUpdateTimer; public SlowEventsLogger( Supplier p99Supplier, - Logger log, - Time time + LogContext logContext ) { this.p99Supplier = p99Supplier; - this.log = log; - this.percentileUpdateTimer = time.timer(P99_REFRESH_INTERVAL_MS); this.p99 = p99Supplier.get(); + this.log = logContext.logger(SlowEventsLogger.class); } - public void maybeLog(String name, long durationNs) { + public void maybeLogEvent(String name, long durationNs) { long durationMs = MILLISECONDS.convert(durationNs, NANOSECONDS); if (durationMs > MIN_SLOW_EVENT_TIME_MS && durationMs > p99) { - log.info("Slow controller event {} processed in {} us", name, - MICROSECONDS.convert(durationNs, NANOSECONDS)); + log.info("Slow controller event {} processed in {} us", + name, MICROSECONDS.convert(durationNs, NANOSECONDS)); } } - public void maybeRefreshPercentile() { - percentileUpdateTimer.update(); - if (percentileUpdateTimer.isExpired()) { - p99 = p99Supplier.get(); - log.trace("Update slow events p99 to {}", p99); - percentileUpdateTimer.reset(P99_REFRESH_INTERVAL_MS); - } + public void refreshPercentile() { + p99 = p99Supplier.get(); + log.trace("Update slow controller event threshold (p99) to {}", p99); } } diff --git a/metadata/src/test/java/org/apache/kafka/controller/metrics/SlowEventsLoggerTest.java b/metadata/src/test/java/org/apache/kafka/controller/metrics/SlowEventsLoggerTest.java new file mode 100644 index 0000000000000..9ec93ab4fc1f4 --- /dev/null +++ b/metadata/src/test/java/org/apache/kafka/controller/metrics/SlowEventsLoggerTest.java @@ -0,0 +1,17 @@ +package org.apache.kafka.controller.metrics; + +import org.apache.kafka.common.utils.LogContext; +import org.junit.jupiter.api.Test; + +import java.util.concurrent.atomic.AtomicReference; + +public class SlowEventsLoggerTest { + @Test + public void test() { + LogContext logContext = new LogContext(); + AtomicReference p99 = new AtomicReference<>(1000.0); + SlowEventsLogger logger = new SlowEventsLogger(p99::get, logContext); + + logger.maybeLogEvent(); + } +} From c20ec567567a77a8480f320bbb085bd508a7fb7f Mon Sep 17 00:00:00 2001 From: David Arthur Date: Thu, 12 Dec 2024 10:24:15 -0500 Subject: [PATCH 4/9] updates --- .../scala/kafka/server/ControllerServer.scala | 3 +- .../main/scala/kafka/server/KafkaConfig.scala | 1 + .../kafka/controller/QuorumController.java | 48 ++++++---- .../kafka/controller/SlowEventsLogger.java | 87 ++++++++++++++++++ .../metrics/QuorumControllerMetrics.java | 6 +- .../controller/metrics/SlowEventsLogger.java | 49 ---------- .../controller/SlowEventsLoggerTest.java | 89 +++++++++++++++++++ .../metrics/SlowEventsLoggerTest.java | 17 ---- .../kafka/server/config/KRaftConfigs.java | 6 +- 9 files changed, 216 insertions(+), 90 deletions(-) create mode 100644 metadata/src/main/java/org/apache/kafka/controller/SlowEventsLogger.java delete mode 100644 metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java create mode 100644 metadata/src/test/java/org/apache/kafka/controller/SlowEventsLoggerTest.java delete mode 100644 metadata/src/test/java/org/apache/kafka/controller/metrics/SlowEventsLoggerTest.java diff --git a/core/src/main/scala/kafka/server/ControllerServer.scala b/core/src/main/scala/kafka/server/ControllerServer.scala index 8cc14516126a0..b3f6e785c942c 100644 --- a/core/src/main/scala/kafka/server/ControllerServer.scala +++ b/core/src/main/scala/kafka/server/ControllerServer.scala @@ -242,7 +242,8 @@ class ControllerServer( setDelegationTokenExpiryTimeMs(config.delegationTokenExpiryTimeMs). setDelegationTokenExpiryCheckIntervalMs(config.delegationTokenExpiryCheckIntervalMs). setUncleanLeaderElectionCheckIntervalMs(config.uncleanLeaderElectionCheckIntervalMs). - setInterBrokerListenerName(config.interBrokerListenerName.value()) + setInterBrokerListenerName(config.interBrokerListenerName.value()). + setMinSlowEventTimeMs(config.minSlowEventTimeMs) } controller = controllerBuilder.build() diff --git a/core/src/main/scala/kafka/server/KafkaConfig.scala b/core/src/main/scala/kafka/server/KafkaConfig.scala index 9502e81f3e49d..a724633a29c19 100755 --- a/core/src/main/scala/kafka/server/KafkaConfig.scala +++ b/core/src/main/scala/kafka/server/KafkaConfig.scala @@ -333,6 +333,7 @@ class KafkaConfig private(doLog: Boolean, val props: util.Map[_, _]) val initialRegistrationTimeoutMs: Int = getInt(KRaftConfigs.INITIAL_BROKER_REGISTRATION_TIMEOUT_MS_CONFIG) val brokerHeartbeatIntervalMs: Int = getInt(KRaftConfigs.BROKER_HEARTBEAT_INTERVAL_MS_CONFIG) val brokerSessionTimeoutMs: Int = getInt(KRaftConfigs.BROKER_SESSION_TIMEOUT_MS_CONFIG) + val minSlowEventTimeMs: Int = getInt(KRaftConfigs.MIN_SLOW_EVENT_TIME_MS_CONFIG) def requiresZookeeper: Boolean = processRoles.isEmpty def usesSelfManagedQuorum: Boolean = processRoles.nonEmpty diff --git a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java index 9c3c934a75a50..cd820692a8cb0 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -94,7 +94,6 @@ import org.apache.kafka.controller.errors.ControllerExceptions; import org.apache.kafka.controller.errors.EventHandlerExceptionInfo; import org.apache.kafka.controller.metrics.QuorumControllerMetrics; -import org.apache.kafka.controller.metrics.SlowEventsLogger; import org.apache.kafka.deferred.DeferredEvent; import org.apache.kafka.deferred.DeferredEventQueue; import org.apache.kafka.metadata.BrokerHeartbeatReply; @@ -152,7 +151,6 @@ import static java.util.concurrent.TimeUnit.MICROSECONDS; import static java.util.concurrent.TimeUnit.MILLISECONDS; import static java.util.concurrent.TimeUnit.NANOSECONDS; -import static java.util.concurrent.TimeUnit.SECONDS; import static org.apache.kafka.controller.QuorumController.ControllerOperationFlag.DOES_NOT_UPDATE_QUEUE_TIME; @@ -178,16 +176,21 @@ */ public final class QuorumController implements Controller { /** - * The maximum records that the controller will write in a single batch. + * The default maximum records that the controller will write in a single batch. */ - private static final int MAX_RECORDS_PER_BATCH = 10000; + private static final int DEFAULT_MAX_RECORDS_PER_BATCH = 10000; + + /** + * The default minimum event time that can be logged as a slow event. + */ + private static final int DEFAULT_MIN_SLOW_EVENT_TIME_MS = 200; /** * The maximum records any user-initiated operation is allowed to generate. * * For now, this is set to the maximum records in a single batch. */ - static final int MAX_RECORDS_PER_USER_OP = MAX_RECORDS_PER_BATCH; + static final int MAX_RECORDS_PER_USER_OP = DEFAULT_MAX_RECORDS_PER_BATCH; /** * A builder class which creates the QuorumController. @@ -216,7 +219,8 @@ public static class Builder { private ConfigurationValidator configurationValidator = ConfigurationValidator.NO_OP; private Map staticConfig = Collections.emptyMap(); private BootstrapMetadata bootstrapMetadata = null; - private int maxRecordsPerBatch = MAX_RECORDS_PER_BATCH; + private int maxRecordsPerBatch = DEFAULT_MAX_RECORDS_PER_BATCH; + private int minSlowEventTimeMs = DEFAULT_MIN_SLOW_EVENT_TIME_MS; private DelegationTokenCache tokenCache; private String tokenSecretKeyString; private long delegationTokenMaxLifeMs; @@ -324,6 +328,11 @@ public Builder setMaxRecordsPerBatch(int maxRecordsPerBatch) { return this; } + public Builder setMinSlowEventTimeMs(int minSlowEventTimeMs) { + this.minSlowEventTimeMs = minSlowEventTimeMs; + return this; + } + public Builder setCreateTopicPolicy(Optional createTopicPolicy) { this.createTopicPolicy = createTopicPolicy; return this; @@ -436,7 +445,8 @@ public QuorumController build() throws Exception { delegationTokenExpiryTimeMs, delegationTokenExpiryCheckIntervalMs, uncleanLeaderElectionCheckIntervalMs, - interBrokerListenerName + interBrokerListenerName, + minSlowEventTimeMs ); } catch (Exception e) { Utils.closeQuietly(queue, "event queue"); @@ -1482,7 +1492,8 @@ private QuorumController( long delegationTokenExpiryTimeMs, long delegationTokenExpiryCheckIntervalMs, long uncleanLeaderElectionCheckIntervalMs, - String interBrokerListenerName + String interBrokerListenerName, + int minSlowEventTimeMs ) { this.nonFatalFaultHandler = nonFatalFaultHandler; this.fatalFaultHandler = fatalFaultHandler; @@ -1592,7 +1603,7 @@ private QuorumController( } registerElectUnclean(TimeUnit.MILLISECONDS.toNanos(uncleanLeaderElectionCheckIntervalMs)); registerExpireDelegationTokens(MILLISECONDS.toNanos(delegationTokenExpiryCheckIntervalMs)); - registerSlowEventUpdater(MILLISECONDS.toNanos(30000)); + registerUpdateSlowEventLogger(MILLISECONDS.toNanos(30000)); // OffsetControlManager must be initialized last, because its constructor will take the // initial in-memory snapshot of all extant timeline data structures. this.offsetControl = new OffsetControlManager.Builder(). @@ -1604,7 +1615,8 @@ private QuorumController( log.info("Creating new QuorumController with clusterId {}", clusterId); this.raftClient.register(metaLogListener); - this.slowEventsLogger = new SlowEventsLogger(controllerMetrics::getEventQueueProcessingTimeP99, logContext); + this.slowEventsLogger = new SlowEventsLogger(minSlowEventTimeMs, + controllerMetrics::getEventQueueProcessingTime99, logContext); } /** @@ -1628,14 +1640,14 @@ private void registerWriteNoOpRecord(long maxIdleIntervalNs) { EnumSet.noneOf(PeriodicTaskFlag.class))); } - private void registerSlowEventUpdater(long maxSlowEventWindowNs) { - periodicControl.registerTask(new PeriodicTask("updateSlowEventP99", - () -> { - slowEventsLogger.refreshPercentile(); - return ControllerResult.of(Collections.emptyList(), false); - }, - maxSlowEventWindowNs, - EnumSet.noneOf(PeriodicTaskFlag.class))); + private void registerUpdateSlowEventLogger(long maxSlowEventWindowNs) { + periodicControl.registerTask(new PeriodicTask("updateSlowEventLoggerP99", + () -> { + slowEventsLogger.refreshPercentile(); + return ControllerResult.of(Collections.emptyList(), false); + }, + maxSlowEventWindowNs, + EnumSet.noneOf(PeriodicTaskFlag.class))); } /** diff --git a/metadata/src/main/java/org/apache/kafka/controller/SlowEventsLogger.java b/metadata/src/main/java/org/apache/kafka/controller/SlowEventsLogger.java new file mode 100644 index 0000000000000..913f08a4db933 --- /dev/null +++ b/metadata/src/main/java/org/apache/kafka/controller/SlowEventsLogger.java @@ -0,0 +1,87 @@ +/* + * 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.controller; + +import org.apache.kafka.common.utils.LogContext; + +import org.slf4j.Logger; + +import java.util.function.Supplier; + +import static java.util.concurrent.TimeUnit.MILLISECONDS; +import static java.util.concurrent.TimeUnit.NANOSECONDS; + +/** + * Track the p99 for controller event queue processing time. If we encounter an event that takes longer + * than this cached p99 time, we will log it at INFO level on the controller logger. + */ +public class SlowEventsLogger { + /** + * Don't report any p99 events below this threshold. This prevents the controller from reporting p99 event + * times in the idle case where p99 event times are essentially the average as well. + */ + private final long minSlowEventTimeNs; + + /** + * Function that returns the current p99 time in millis. This call can be expensive, and since the histogram is + * biased towards the last 5 minutes of data, we only need to update this p99 every so often. + */ + private final Supplier thresholdMsSupplier; + + private final Logger log; + + /** + * The current p99 threshold in nanos. + */ + private long thresholdNs; + + public SlowEventsLogger( + int minSlowEventTimeMs, + Supplier thresholdMsSupplier, + LogContext logContext + ) { + this.minSlowEventTimeNs = MILLISECONDS.toNanos(minSlowEventTimeMs); + this.thresholdMsSupplier = thresholdMsSupplier; + this.thresholdNs = minSlowEventTimeMs; + this.log = logContext.logger(SlowEventsLogger.class); + } + + /** + * Produce an INFO log if the given event ran for longer than the current p99 event processing time. + * + * @return true if a slow event was logged, false otherwise. + */ + public boolean maybeLogEvent(String name, long durationNs) { + if (durationNs >= minSlowEventTimeNs && durationNs >= thresholdNs) { + log.info("Slow controller event {} processed in {} us which is greater than p99 of {} us", + name, + NANOSECONDS.toMicros(durationNs), + NANOSECONDS.toMicros(thresholdNs) + ); + return true; + } else if (durationNs >= thresholdNs) { + return false; + } + return false; + } + + public void refreshPercentile() { + thresholdNs = (long) (thresholdMsSupplier.get() * 1000000); + log.trace("Update slow controller event threshold (p99) to {} us.", NANOSECONDS.toMicros(thresholdNs)); + } +} diff --git a/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java b/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java index ada68a4ff210c..e27bd3b83d9e7 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java +++ b/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java @@ -29,8 +29,6 @@ import java.util.Optional; import java.util.concurrent.atomic.AtomicLong; import java.util.function.Consumer; -import java.util.function.Function; -import java.util.function.Supplier; /** * These are the metrics which are managed by the QuorumController class. They generally pertain to @@ -166,9 +164,9 @@ public void updateEventQueueProcessingTime(long durationMs) { eventQueueProcessingTimeUpdater.accept(durationMs); } - public double getEventQueueProcessingTimeP99() { + public double getEventQueueProcessingTime99() { if (registry.isPresent()) { - Histogram histogram = registry.get().newHistogram(EVENT_QUEUE_PROCESSING_TIME_MS, false); + Histogram histogram = registry.get().newHistogram(EVENT_QUEUE_PROCESSING_TIME_MS, true); return histogram.getSnapshot().get99thPercentile(); } else { // Only returned in unit tests when a metrics registry is not set. diff --git a/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java b/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java deleted file mode 100644 index fb2a128e50d3e..0000000000000 --- a/metadata/src/main/java/org/apache/kafka/controller/metrics/SlowEventsLogger.java +++ /dev/null @@ -1,49 +0,0 @@ -package org.apache.kafka.controller.metrics; - -import org.apache.kafka.common.utils.LogContext; -import org.apache.kafka.controller.QuorumController; -import org.slf4j.Logger; - -import java.util.function.Supplier; - -import static java.util.concurrent.TimeUnit.MICROSECONDS; -import static java.util.concurrent.TimeUnit.MILLISECONDS; -import static java.util.concurrent.TimeUnit.NANOSECONDS; - -/** - * Track the p99 for controller event queue processing time. If we encounter an event that takes longer - * than this cached p99 time, we will log it at INFO level on the controller logger. - */ -public class SlowEventsLogger { - /** - * Don't report any p99 events below this threshold. This prevents the controller from reporting p99 event - * times in the idle case. - */ - private static final int MIN_SLOW_EVENT_TIME_MS = 100; - - private final Supplier p99Supplier; - private final Logger log; - private double p99; - - public SlowEventsLogger( - Supplier p99Supplier, - LogContext logContext - ) { - this.p99Supplier = p99Supplier; - this.p99 = p99Supplier.get(); - this.log = logContext.logger(SlowEventsLogger.class); - } - - public void maybeLogEvent(String name, long durationNs) { - long durationMs = MILLISECONDS.convert(durationNs, NANOSECONDS); - if (durationMs > MIN_SLOW_EVENT_TIME_MS && durationMs > p99) { - log.info("Slow controller event {} processed in {} us", - name, MICROSECONDS.convert(durationNs, NANOSECONDS)); - } - } - - public void refreshPercentile() { - p99 = p99Supplier.get(); - log.trace("Update slow controller event threshold (p99) to {}", p99); - } -} diff --git a/metadata/src/test/java/org/apache/kafka/controller/SlowEventsLoggerTest.java b/metadata/src/test/java/org/apache/kafka/controller/SlowEventsLoggerTest.java new file mode 100644 index 0000000000000..0f423cce78fe4 --- /dev/null +++ b/metadata/src/test/java/org/apache/kafka/controller/SlowEventsLoggerTest.java @@ -0,0 +1,89 @@ +/* + * 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.controller; + +import org.apache.kafka.common.utils.LogContext; + +import org.junit.jupiter.api.Test; + +import java.util.concurrent.atomic.AtomicReference; + +import static java.util.concurrent.TimeUnit.MILLISECONDS; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + + +public class SlowEventsLoggerTest { + @Test + public void testSlowEvents() { + LogContext logContext = new LogContext(); + + AtomicReference p99 = new AtomicReference<>(0.0); + SlowEventsLogger logger = new SlowEventsLogger(100, p99::get, logContext); + + // Initially, the p99 is zero + assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(10))); + assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(99))); + assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(100))); + + + // Idle controller, low p99 + p99.set(30.0); + logger.refreshPercentile(); + assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(90))); + assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(99))); + assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(100))); + + // Busy controller, high p99 + p99.set(1000.0); + logger.refreshPercentile(); + assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(100))); + assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(200))); + assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(1000))); + assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(2000))); + } + + @Test + public void testThresholdDisabled() { + LogContext logContext = new LogContext(); + + AtomicReference p99 = new AtomicReference<>(0.0); + // Set min slow event time to zero, effectively disabling the threshold + SlowEventsLogger logger = new SlowEventsLogger(0, p99::get, logContext); + + assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(0))); + assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(10))); + assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(99))); + assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(100))); + + p99.set(30.0); + logger.refreshPercentile(); + assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(0))); + assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(29))); + assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(30))); + assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(100))); + + + p99.set(1000.0); + logger.refreshPercentile(); + assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(100))); + assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(999))); + assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(1000))); + assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(2000))); + } +} diff --git a/metadata/src/test/java/org/apache/kafka/controller/metrics/SlowEventsLoggerTest.java b/metadata/src/test/java/org/apache/kafka/controller/metrics/SlowEventsLoggerTest.java deleted file mode 100644 index 9ec93ab4fc1f4..0000000000000 --- a/metadata/src/test/java/org/apache/kafka/controller/metrics/SlowEventsLoggerTest.java +++ /dev/null @@ -1,17 +0,0 @@ -package org.apache.kafka.controller.metrics; - -import org.apache.kafka.common.utils.LogContext; -import org.junit.jupiter.api.Test; - -import java.util.concurrent.atomic.AtomicReference; - -public class SlowEventsLoggerTest { - @Test - public void test() { - LogContext logContext = new LogContext(); - AtomicReference p99 = new AtomicReference<>(1000.0); - SlowEventsLogger logger = new SlowEventsLogger(p99::get, logContext); - - logger.maybeLogEvent(); - } -} diff --git a/server/src/main/java/org/apache/kafka/server/config/KRaftConfigs.java b/server/src/main/java/org/apache/kafka/server/config/KRaftConfigs.java index d2cf35f33a5bc..866c79d6f7c66 100644 --- a/server/src/main/java/org/apache/kafka/server/config/KRaftConfigs.java +++ b/server/src/main/java/org/apache/kafka/server/config/KRaftConfigs.java @@ -52,7 +52,6 @@ public class KRaftConfigs { public static final int BROKER_SESSION_TIMEOUT_MS_DEFAULT = 9000; public static final String BROKER_SESSION_TIMEOUT_MS_DOC = "The length of time in milliseconds that a broker lease lasts if no heartbeats are made. Used when running in KRaft mode."; - public static final String NODE_ID_CONFIG = "node.id"; public static final int EMPTY_NODE_ID = -1; public static final String NODE_ID_DOC = "The node ID associated with the roles this process is playing when process.roles is non-empty. " + @@ -115,6 +114,10 @@ public class KRaftConfigs { public static final String SERVER_MAX_STARTUP_TIME_MS_DOC = "The maximum number of milliseconds we will wait for the server to come up. " + "By default there is no limit. This should be used for testing only."; + public static final String MIN_SLOW_EVENT_TIME_MS_CONFIG = "controller.slow.event.min.ms"; + public static final int MIN_SLOW_EVENT_TIME_MS_DEFAULT = 200; + public static final String MIN_SLOW_EVENT_TIME_MS_DOC = "Do not log slow controller events which are slowed than this duration."; + /** ZK to KRaft Migration configs */ public static final String MIGRATION_ENABLED_CONFIG = "zookeeper.metadata.migration.enable"; public static final String MIGRATION_ENABLED_DOC = "Enable ZK to KRaft migration"; @@ -141,6 +144,7 @@ public class KRaftConfigs { .define(METADATA_MAX_RETENTION_MILLIS_CONFIG, LONG, LogConfig.DEFAULT_RETENTION_MS, null, HIGH, METADATA_MAX_RETENTION_MILLIS_DOC) .define(METADATA_MAX_IDLE_INTERVAL_MS_CONFIG, INT, METADATA_MAX_IDLE_INTERVAL_MS_DEFAULT, atLeast(0), LOW, METADATA_MAX_IDLE_INTERVAL_MS_DOC) .defineInternal(SERVER_MAX_STARTUP_TIME_MS_CONFIG, LONG, SERVER_MAX_STARTUP_TIME_MS_DEFAULT, atLeast(0), MEDIUM, SERVER_MAX_STARTUP_TIME_MS_DOC) + .defineInternal(MIN_SLOW_EVENT_TIME_MS_CONFIG, INT, MIN_SLOW_EVENT_TIME_MS_DEFAULT, atLeast(0), MEDIUM, MIN_SLOW_EVENT_TIME_MS_DOC) .define(MIGRATION_ENABLED_CONFIG, BOOLEAN, false, HIGH, MIGRATION_ENABLED_DOC) .defineInternal(MIGRATION_METADATA_MIN_BATCH_SIZE_CONFIG, INT, MIGRATION_METADATA_MIN_BATCH_SIZE_DEFAULT, atLeast(1), MEDIUM, MIGRATION_METADATA_MIN_BATCH_SIZE_DOC); From c010cb4a7c363fcb5192a637c998c348667e133e Mon Sep 17 00:00:00 2001 From: David Arthur Date: Thu, 12 Dec 2024 20:40:01 -0500 Subject: [PATCH 5/9] pr feedback --- .../java/org/apache/kafka/controller/QuorumController.java | 3 ++- .../java/org/apache/kafka/controller/SlowEventsLogger.java | 6 ++---- .../java/org/apache/kafka/server/config/KRaftConfigs.java | 2 +- 3 files changed, 5 insertions(+), 6 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java index 2db26dd671fe7..df0a42654a481 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -150,6 +150,7 @@ import static java.util.concurrent.TimeUnit.MICROSECONDS; import static java.util.concurrent.TimeUnit.MILLISECONDS; import static java.util.concurrent.TimeUnit.NANOSECONDS; +import static java.util.concurrent.TimeUnit.SECONDS; import static org.apache.kafka.controller.QuorumController.ControllerOperationFlag.DOES_NOT_UPDATE_QUEUE_TIME; @@ -1603,7 +1604,7 @@ private QuorumController( } registerElectUnclean(TimeUnit.MILLISECONDS.toNanos(uncleanLeaderElectionCheckIntervalMs)); registerExpireDelegationTokens(MILLISECONDS.toNanos(delegationTokenExpiryCheckIntervalMs)); - registerUpdateSlowEventLogger(MILLISECONDS.toNanos(30000)); + registerUpdateSlowEventLogger(SECONDS.toNanos(30)); // OffsetControlManager must be initialized last, because its constructor will take the // initial in-memory snapshot of all extant timeline data structures. this.offsetControl = new OffsetControlManager.Builder(). diff --git a/metadata/src/main/java/org/apache/kafka/controller/SlowEventsLogger.java b/metadata/src/main/java/org/apache/kafka/controller/SlowEventsLogger.java index 913f08a4db933..691e2ae313f82 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/SlowEventsLogger.java +++ b/metadata/src/main/java/org/apache/kafka/controller/SlowEventsLogger.java @@ -62,20 +62,18 @@ public SlowEventsLogger( } /** - * Produce an INFO log if the given event ran for longer than the current p99 event processing time. + * Produce an INFO log if the given event ran for at least as long as the current p99 event processing time. * * @return true if a slow event was logged, false otherwise. */ public boolean maybeLogEvent(String name, long durationNs) { if (durationNs >= minSlowEventTimeNs && durationNs >= thresholdNs) { - log.info("Slow controller event {} processed in {} us which is greater than p99 of {} us", + log.info("Slow controller event {} processed in {} us which is larger or equal to the p99 of {} us", name, NANOSECONDS.toMicros(durationNs), NANOSECONDS.toMicros(thresholdNs) ); return true; - } else if (durationNs >= thresholdNs) { - return false; } return false; } diff --git a/server/src/main/java/org/apache/kafka/server/config/KRaftConfigs.java b/server/src/main/java/org/apache/kafka/server/config/KRaftConfigs.java index 4787509497aa3..6ad4647796416 100644 --- a/server/src/main/java/org/apache/kafka/server/config/KRaftConfigs.java +++ b/server/src/main/java/org/apache/kafka/server/config/KRaftConfigs.java @@ -116,7 +116,7 @@ public class KRaftConfigs { public static final String MIN_SLOW_EVENT_TIME_MS_CONFIG = "controller.slow.event.min.ms"; public static final int MIN_SLOW_EVENT_TIME_MS_DEFAULT = 200; - public static final String MIN_SLOW_EVENT_TIME_MS_DOC = "Do not log slow controller events which are slowed than this duration."; + public static final String MIN_SLOW_EVENT_TIME_MS_DOC = "Log controller events with a p99 duration slower than this amount."; public static final ConfigDef CONFIG_DEF = new ConfigDef() .define(METADATA_SNAPSHOT_MAX_NEW_RECORD_BYTES_CONFIG, LONG, METADATA_SNAPSHOT_MAX_NEW_RECORD_BYTES, atLeast(1), HIGH, METADATA_SNAPSHOT_MAX_NEW_RECORD_BYTES_DOC) From d23e4af0b8213c1d295d1015b33fd03662d596c4 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Fri, 13 Dec 2024 21:28:50 -0500 Subject: [PATCH 6/9] pr feedback --- ...gger.java => EventPerformanceMonitor.java} | 31 ++++++++++-- .../kafka/controller/QuorumController.java | 8 +-- ....java => EventPerformanceMonitorTest.java} | 50 +++++++++---------- 3 files changed, 55 insertions(+), 34 deletions(-) rename metadata/src/main/java/org/apache/kafka/controller/{SlowEventsLogger.java => EventPerformanceMonitor.java} (76%) rename metadata/src/test/java/org/apache/kafka/controller/{SlowEventsLoggerTest.java => EventPerformanceMonitorTest.java} (52%) diff --git a/metadata/src/main/java/org/apache/kafka/controller/SlowEventsLogger.java b/metadata/src/main/java/org/apache/kafka/controller/EventPerformanceMonitor.java similarity index 76% rename from metadata/src/main/java/org/apache/kafka/controller/SlowEventsLogger.java rename to metadata/src/main/java/org/apache/kafka/controller/EventPerformanceMonitor.java index 691e2ae313f82..b80eab7a15086 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/SlowEventsLogger.java +++ b/metadata/src/main/java/org/apache/kafka/controller/EventPerformanceMonitor.java @@ -30,7 +30,7 @@ * Track the p99 for controller event queue processing time. If we encounter an event that takes longer * than this cached p99 time, we will log it at INFO level on the controller logger. */ -public class SlowEventsLogger { +public class EventPerformanceMonitor { /** * Don't report any p99 events below this threshold. This prevents the controller from reporting p99 event * times in the idle case where p99 event times are essentially the average as well. @@ -50,7 +50,11 @@ public class SlowEventsLogger { */ private long thresholdNs; - public SlowEventsLogger( + private String slowestEvent; + private long slowestDurationNs; + + + public EventPerformanceMonitor( int minSlowEventTimeMs, Supplier thresholdMsSupplier, LogContext logContext @@ -58,7 +62,9 @@ public SlowEventsLogger( this.minSlowEventTimeNs = MILLISECONDS.toNanos(minSlowEventTimeMs); this.thresholdMsSupplier = thresholdMsSupplier; this.thresholdNs = minSlowEventTimeMs; - this.log = logContext.logger(SlowEventsLogger.class); + this.log = logContext.logger(EventPerformanceMonitor.class); + this.slowestEvent = null; + this.slowestDurationNs = 0; } /** @@ -66,8 +72,17 @@ public SlowEventsLogger( * * @return true if a slow event was logged, false otherwise. */ - public boolean maybeLogEvent(String name, long durationNs) { - if (durationNs >= minSlowEventTimeNs && durationNs >= thresholdNs) { + public boolean observeEvent(String name, long durationNs) { + if (durationNs < minSlowEventTimeNs) { + return false; + } + + if (slowestEvent == null || slowestDurationNs < durationNs) { + slowestEvent = name; + slowestDurationNs = durationNs; + } + + if (durationNs >= thresholdNs) { log.info("Slow controller event {} processed in {} us which is larger or equal to the p99 of {} us", name, NANOSECONDS.toMicros(durationNs), @@ -79,6 +94,12 @@ public boolean maybeLogEvent(String name, long durationNs) { } public void refreshPercentile() { + if (slowestEvent != null) { + log.info("Slowest event since last refresh was {} which processed in {} us", + slowestEvent, NANOSECONDS.toMicros(slowestDurationNs)); + slowestEvent = null; + slowestDurationNs = 0; + } thresholdNs = (long) (thresholdMsSupplier.get() * 1000000); log.trace("Update slow controller event threshold (p99) to {} us.", NANOSECONDS.toMicros(thresholdNs)); } diff --git a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java index df0a42654a481..25ec33fffe4fd 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -537,7 +537,7 @@ private void handleEventEnd(String name, long startProcessingTimeNs) { long deltaNs = endProcessingTime - startProcessingTimeNs; log.debug("Processed {} in {} us", name, MICROSECONDS.convert(deltaNs, NANOSECONDS)); - slowEventsLogger.maybeLogEvent(name, deltaNs); + eventPerformanceMonitor.observeEvent(name, deltaNs); controllerMetrics.updateEventQueueProcessingTime(NANOSECONDS.toMillis(deltaNs)); } @@ -1460,7 +1460,7 @@ private void replay(ApiMessage message, Optional snapshotId, lon */ private final RecordRedactor recordRedactor; - private final SlowEventsLogger slowEventsLogger; + private final EventPerformanceMonitor eventPerformanceMonitor; private QuorumController( FaultHandler nonFatalFaultHandler, @@ -1616,7 +1616,7 @@ private QuorumController( log.info("Creating new QuorumController with clusterId {}", clusterId); this.raftClient.register(metaLogListener); - this.slowEventsLogger = new SlowEventsLogger(minSlowEventTimeMs, + this.eventPerformanceMonitor = new EventPerformanceMonitor(minSlowEventTimeMs, controllerMetrics::getEventQueueProcessingTime99, logContext); } @@ -1644,7 +1644,7 @@ private void registerWriteNoOpRecord(long maxIdleIntervalNs) { private void registerUpdateSlowEventLogger(long maxSlowEventWindowNs) { periodicControl.registerTask(new PeriodicTask("updateSlowEventLoggerP99", () -> { - slowEventsLogger.refreshPercentile(); + eventPerformanceMonitor.refreshPercentile(); return ControllerResult.of(Collections.emptyList(), false); }, maxSlowEventWindowNs, diff --git a/metadata/src/test/java/org/apache/kafka/controller/SlowEventsLoggerTest.java b/metadata/src/test/java/org/apache/kafka/controller/EventPerformanceMonitorTest.java similarity index 52% rename from metadata/src/test/java/org/apache/kafka/controller/SlowEventsLoggerTest.java rename to metadata/src/test/java/org/apache/kafka/controller/EventPerformanceMonitorTest.java index 0f423cce78fe4..814f4f36a900b 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/SlowEventsLoggerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/EventPerformanceMonitorTest.java @@ -28,34 +28,34 @@ import static org.junit.jupiter.api.Assertions.assertTrue; -public class SlowEventsLoggerTest { +public class EventPerformanceMonitorTest { @Test public void testSlowEvents() { LogContext logContext = new LogContext(); AtomicReference p99 = new AtomicReference<>(0.0); - SlowEventsLogger logger = new SlowEventsLogger(100, p99::get, logContext); + EventPerformanceMonitor logger = new EventPerformanceMonitor(100, p99::get, logContext); // Initially, the p99 is zero - assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(10))); - assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(99))); - assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(100))); + assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(10))); + assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(99))); + assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(100))); // Idle controller, low p99 p99.set(30.0); logger.refreshPercentile(); - assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(90))); - assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(99))); - assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(100))); + assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(90))); + assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(99))); + assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(100))); // Busy controller, high p99 p99.set(1000.0); logger.refreshPercentile(); - assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(100))); - assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(200))); - assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(1000))); - assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(2000))); + assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(100))); + assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(200))); + assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(1000))); + assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(2000))); } @Test @@ -64,26 +64,26 @@ public void testThresholdDisabled() { AtomicReference p99 = new AtomicReference<>(0.0); // Set min slow event time to zero, effectively disabling the threshold - SlowEventsLogger logger = new SlowEventsLogger(0, p99::get, logContext); + EventPerformanceMonitor logger = new EventPerformanceMonitor(0, p99::get, logContext); - assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(0))); - assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(10))); - assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(99))); - assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(100))); + assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(0))); + assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(10))); + assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(99))); + assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(100))); p99.set(30.0); logger.refreshPercentile(); - assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(0))); - assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(29))); - assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(30))); - assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(100))); + assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(0))); + assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(29))); + assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(30))); + assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(100))); p99.set(1000.0); logger.refreshPercentile(); - assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(100))); - assertFalse(logger.maybeLogEvent("test", MILLISECONDS.toNanos(999))); - assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(1000))); - assertTrue(logger.maybeLogEvent("test", MILLISECONDS.toNanos(2000))); + assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(100))); + assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(999))); + assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(1000))); + assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(2000))); } } From a2ded76accec83e9e3708892d22a88c77a83aaf8 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Sun, 15 Dec 2024 21:45:12 -0500 Subject: [PATCH 7/9] set eventPerformanceMonitor earlier --- .../java/org/apache/kafka/controller/QuorumController.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java index 25ec33fffe4fd..865324ae85765 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -1591,6 +1591,8 @@ private QuorumController( this.metaLogListener = new QuorumMetaLogListener(); this.curClaimEpoch = -1; this.recordRedactor = new RecordRedactor(configSchema); + this.eventPerformanceMonitor = new EventPerformanceMonitor(minSlowEventTimeMs, + controllerMetrics::getEventQueueProcessingTime99, logContext); if (maxIdleIntervalNs.isPresent()) { registerWriteNoOpRecord(maxIdleIntervalNs.getAsLong()); } @@ -1614,10 +1616,7 @@ private QuorumController( setTime(time). build(); log.info("Creating new QuorumController with clusterId {}", clusterId); - this.raftClient.register(metaLogListener); - this.eventPerformanceMonitor = new EventPerformanceMonitor(minSlowEventTimeMs, - controllerMetrics::getEventQueueProcessingTime99, logContext); } /** From 8993f1056296c1165ae9f0331fd42e4b5f891e44 Mon Sep 17 00:00:00 2001 From: "Colin P. McCabe" Date: Fri, 20 Dec 2024 11:23:24 -0800 Subject: [PATCH 8/9] logging changes etc. --- .../scala/kafka/server/ControllerServer.scala | 3 +- .../main/scala/kafka/server/KafkaConfig.scala | 3 +- .../controller/EventPerformanceMonitor.java | 206 +++++++++++++----- .../kafka/controller/QuorumController.java | 62 ++++-- .../metrics/QuorumControllerMetrics.java | 9 - .../EventPerformanceMonitorTest.java | 147 ++++++++----- .../kafka/server/config/KRaftConfigs.java | 13 +- 7 files changed, 299 insertions(+), 144 deletions(-) diff --git a/core/src/main/scala/kafka/server/ControllerServer.scala b/core/src/main/scala/kafka/server/ControllerServer.scala index b3f6e785c942c..db2631ef149a8 100644 --- a/core/src/main/scala/kafka/server/ControllerServer.scala +++ b/core/src/main/scala/kafka/server/ControllerServer.scala @@ -243,7 +243,8 @@ class ControllerServer( setDelegationTokenExpiryCheckIntervalMs(config.delegationTokenExpiryCheckIntervalMs). setUncleanLeaderElectionCheckIntervalMs(config.uncleanLeaderElectionCheckIntervalMs). setInterBrokerListenerName(config.interBrokerListenerName.value()). - setMinSlowEventTimeMs(config.minSlowEventTimeMs) + setControllerPerformanceSamplePeriodMs(config.controllerPerformanceSamplePeriodMs). + setControllerPerformanceAlwaysLogThresholdMs(config.controllerPerformanceAlwaysLogThresholdMs) } controller = controllerBuilder.build() diff --git a/core/src/main/scala/kafka/server/KafkaConfig.scala b/core/src/main/scala/kafka/server/KafkaConfig.scala index 916277ed07ed7..5028c8c67ff46 100755 --- a/core/src/main/scala/kafka/server/KafkaConfig.scala +++ b/core/src/main/scala/kafka/server/KafkaConfig.scala @@ -333,7 +333,8 @@ class KafkaConfig private(doLog: Boolean, val props: util.Map[_, _]) val initialRegistrationTimeoutMs: Int = getInt(KRaftConfigs.INITIAL_BROKER_REGISTRATION_TIMEOUT_MS_CONFIG) val brokerHeartbeatIntervalMs: Int = getInt(KRaftConfigs.BROKER_HEARTBEAT_INTERVAL_MS_CONFIG) val brokerSessionTimeoutMs: Int = getInt(KRaftConfigs.BROKER_SESSION_TIMEOUT_MS_CONFIG) - val minSlowEventTimeMs: Int = getInt(KRaftConfigs.MIN_SLOW_EVENT_TIME_MS_CONFIG) + val controllerPerformanceSamplePeriodMs: Long = getLong(KRaftConfigs.CONTROLLER_PERFORMANCE_SAMPLE_PERIOD_MS) + val controllerPerformanceAlwaysLogThresholdMs: Long = getLong(KRaftConfigs.CONTROLLER_PERFORMANCE_ALWAYS_LOG_THRESHOLD_MS) def requiresZookeeper: Boolean = processRoles.isEmpty def usesSelfManagedQuorum: Boolean = processRoles.nonEmpty diff --git a/metadata/src/main/java/org/apache/kafka/controller/EventPerformanceMonitor.java b/metadata/src/main/java/org/apache/kafka/controller/EventPerformanceMonitor.java index b80eab7a15086..5f0b2635122ee 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/EventPerformanceMonitor.java +++ b/metadata/src/main/java/org/apache/kafka/controller/EventPerformanceMonitor.java @@ -21,86 +21,192 @@ import org.slf4j.Logger; -import java.util.function.Supplier; +import java.text.DecimalFormat; +import java.util.AbstractMap; +import java.util.Map; -import static java.util.concurrent.TimeUnit.MILLISECONDS; import static java.util.concurrent.TimeUnit.NANOSECONDS; +import static java.util.concurrent.TimeUnit.SECONDS; /** - * Track the p99 for controller event queue processing time. If we encounter an event that takes longer - * than this cached p99 time, we will log it at INFO level on the controller logger. + * Track the performance of controller events. Periodically log the slowest events. + * Log any event slower than a certain threshold. */ -public class EventPerformanceMonitor { +class EventPerformanceMonitor { /** - * Don't report any p99 events below this threshold. This prevents the controller from reporting p99 event - * times in the idle case where p99 event times are essentially the average as well. + * The format to use when displaying milliseconds. */ - private final long minSlowEventTimeNs; + private static final DecimalFormat MILLISECOND_DECIMAL_FORMAT = new DecimalFormat("#0.00"); + + static class Builder { + LogContext logContext = null; + long periodNs = SECONDS.toNanos(60); + long alwaysLogThresholdNs = SECONDS.toNanos(2); + + Builder setLogContext(LogContext logContext) { + this.logContext = logContext; + return this; + } + + Builder setPeriodNs(long periodNs) { + this.periodNs = periodNs; + return this; + } + + Builder setAlwaysLogThresholdNs(long alwaysLogThresholdNs) { + this.alwaysLogThresholdNs = alwaysLogThresholdNs; + return this; + } + + EventPerformanceMonitor build() { + if (logContext == null) logContext = new LogContext(); + return new EventPerformanceMonitor(logContext, + periodNs, + alwaysLogThresholdNs); + } + } /** - * Function that returns the current p99 time in millis. This call can be expensive, and since the histogram is - * biased towards the last 5 minutes of data, we only need to update this p99 every so often. + * The log4j object to use. */ - private final Supplier thresholdMsSupplier; - private final Logger log; /** - * The current p99 threshold in nanos. + * The period in nanoseconds. */ - private long thresholdNs; + private long periodNs; - private String slowestEvent; - private long slowestDurationNs; + /** + * The always-log threshold in nanoseconds. + */ + private long alwaysLogThresholdNs; + /** + * The name of the slowest event we've seen so far, or null if none has been seen. + */ + private String slowestEventName; + + /** + * The duration of the slowest event we've seen so far, or 0 if none has been seen. + */ + private long slowestEventDurationNs; + + /** + * The total duration of all the events we've seen. + */ + private long totalEventDurationNs; + + /** + * The number of events we've seen. + */ + private int numEvents; - public EventPerformanceMonitor( - int minSlowEventTimeMs, - Supplier thresholdMsSupplier, - LogContext logContext + private EventPerformanceMonitor( + LogContext logContext, + long periodNs, + long alwaysLogThresholdNs ) { - this.minSlowEventTimeNs = MILLISECONDS.toNanos(minSlowEventTimeMs); - this.thresholdMsSupplier = thresholdMsSupplier; - this.thresholdNs = minSlowEventTimeMs; this.log = logContext.logger(EventPerformanceMonitor.class); - this.slowestEvent = null; - this.slowestDurationNs = 0; + this.periodNs = periodNs; + this.alwaysLogThresholdNs = alwaysLogThresholdNs; + reset(); + } + + long periodNs() { + return periodNs; + } + + Map.Entry slowestEvent() { + return new AbstractMap.SimpleImmutableEntry<>(slowestEventName, slowestEventDurationNs); + } + + /** + * Reset all internal state. + */ + void reset() { + this.slowestEventName = null; + this.slowestEventDurationNs = 0; + this.totalEventDurationNs = 0; + this.numEvents = 0; } /** - * Produce an INFO log if the given event ran for at least as long as the current p99 event processing time. + * Handle a controller event being finished. * - * @return true if a slow event was logged, false otherwise. + * @param name The name of the controller event. + * @param durationNs The duration of the controller event in nanoseconds. */ - public boolean observeEvent(String name, long durationNs) { - if (durationNs < minSlowEventTimeNs) { - return false; + void observeEvent(String name, long durationNs) { + String message = doObserveEvent(name, durationNs); + if (message != null) { + log.error("{}", message); } + } - if (slowestEvent == null || slowestDurationNs < durationNs) { - slowestEvent = name; - slowestDurationNs = durationNs; + /** + * Handle a controller event being finished. + * + * @param name The name of the controller event. + * @param durationNs The duration of the controller event in nanoseconds. + * + * @return The message to log, or null otherwise. + */ + String doObserveEvent(String name, long durationNs) { + if (slowestEventName == null || slowestEventDurationNs < durationNs) { + slowestEventName = name; + slowestEventDurationNs = durationNs; } - - if (durationNs >= thresholdNs) { - log.info("Slow controller event {} processed in {} us which is larger or equal to the p99 of {} us", - name, - NANOSECONDS.toMicros(durationNs), - NANOSECONDS.toMicros(thresholdNs) - ); - return true; + totalEventDurationNs += durationNs; + numEvents++; + if (durationNs < alwaysLogThresholdNs) { + return null; } - return false; + return "Exceptionally slow controller event " + name + " took " + + NANOSECONDS.toMillis(durationNs) + " ms."; } - public void refreshPercentile() { - if (slowestEvent != null) { - log.info("Slowest event since last refresh was {} which processed in {} us", - slowestEvent, NANOSECONDS.toMicros(slowestDurationNs)); - slowestEvent = null; - slowestDurationNs = 0; + /** + * Generate a log message summarizing the events of the last period, + * and then reset our internal state. + */ + void generatePeriodicPerformanceMessage() { + String message = periodicPerformanceMessage(); + log.info("{}", message); + reset(); + } + + /** + * Generate a log message summarizing the events of the last period. + * + * @return The summary string. + */ + String periodicPerformanceMessage() { + StringBuilder bld = new StringBuilder(); + bld.append("In the last "); + bld.append(NANOSECONDS.toMillis(periodNs)); + bld.append(" ms period, "); + if (numEvents == 0) { + bld.append("there were no controller events completed."); + } else { + bld.append(numEvents).append(" controller events were completed, which took an average of "); + bld.append(nanosecondsToDecimalMillis(totalEventDurationNs / numEvents)); + bld.append(" ms each. The slowest event was ").append(slowestEventName); + bld.append(", which took "); + bld.append(nanosecondsToDecimalMillis(slowestEventDurationNs)); + bld.append(" ms."); } - thresholdNs = (long) (thresholdMsSupplier.get() * 1000000); - log.trace("Update slow controller event threshold (p99) to {} us.", NANOSECONDS.toMicros(thresholdNs)); + return bld.toString(); + } + + /** + * Translate a duration in nanoseconds to a decimal duration in milliseconds. + * + * @param durationNs The duration in nanoseconds. + * @return The decimal duration in milliseconds. + */ + static String nanosecondsToDecimalMillis(long durationNs) { + double number = NANOSECONDS.toMicros(durationNs); + number /= 1000; + return MILLISECOND_DECIMAL_FORMAT.format(number); } } diff --git a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java index 865324ae85765..acf5c87f8163e 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -150,7 +150,6 @@ import static java.util.concurrent.TimeUnit.MICROSECONDS; import static java.util.concurrent.TimeUnit.MILLISECONDS; import static java.util.concurrent.TimeUnit.NANOSECONDS; -import static java.util.concurrent.TimeUnit.SECONDS; import static org.apache.kafka.controller.QuorumController.ControllerOperationFlag.DOES_NOT_UPDATE_QUEUE_TIME; @@ -220,7 +219,8 @@ public static class Builder { private Map staticConfig = Collections.emptyMap(); private BootstrapMetadata bootstrapMetadata = null; private int maxRecordsPerBatch = DEFAULT_MAX_RECORDS_PER_BATCH; - private int minSlowEventTimeMs = DEFAULT_MIN_SLOW_EVENT_TIME_MS; + private long controllerPerformanceSamplePeriodMs = 60000L; + private long controllerPerformanceAlwaysLogThresholdMs = 2000L; private DelegationTokenCache tokenCache; private String tokenSecretKeyString; private long delegationTokenMaxLifeMs; @@ -328,8 +328,13 @@ public Builder setMaxRecordsPerBatch(int maxRecordsPerBatch) { return this; } - public Builder setMinSlowEventTimeMs(int minSlowEventTimeMs) { - this.minSlowEventTimeMs = minSlowEventTimeMs; + public Builder setControllerPerformanceSamplePeriodMs(long controllerPerformanceSamplePeriodMs) { + this.controllerPerformanceSamplePeriodMs = controllerPerformanceSamplePeriodMs; + return this; + } + + public Builder setControllerPerformanceAlwaysLogThresholdMs(long controllerPerformanceAlwaysLogThresholdMs) { + this.controllerPerformanceAlwaysLogThresholdMs = controllerPerformanceAlwaysLogThresholdMs; return this; } @@ -446,7 +451,8 @@ public QuorumController build() throws Exception { delegationTokenExpiryCheckIntervalMs, uncleanLeaderElectionCheckIntervalMs, interBrokerListenerName, - minSlowEventTimeMs + controllerPerformanceSamplePeriodMs, + controllerPerformanceAlwaysLogThresholdMs ); } catch (Exception e) { Utils.closeQuietly(queue, "event queue"); @@ -537,7 +543,7 @@ private void handleEventEnd(String name, long startProcessingTimeNs) { long deltaNs = endProcessingTime - startProcessingTimeNs; log.debug("Processed {} in {} us", name, MICROSECONDS.convert(deltaNs, NANOSECONDS)); - eventPerformanceMonitor.observeEvent(name, deltaNs); + performanceMonitor.observeEvent(name, deltaNs); controllerMetrics.updateEventQueueProcessingTime(NANOSECONDS.toMillis(deltaNs)); } @@ -550,6 +556,8 @@ private Throwable handleEventException( if (startProcessingTimeNs.isPresent()) { long endProcessingTime = time.nanoseconds(); long deltaNs = endProcessingTime - startProcessingTimeNs.getAsLong(); + performanceMonitor.observeEvent(name, deltaNs); + controllerMetrics.updateEventQueueProcessingTime(NANOSECONDS.toMillis(deltaNs)); deltaUs = OptionalLong.of(MICROSECONDS.convert(deltaNs, NANOSECONDS)); } else { deltaUs = OptionalLong.empty(); @@ -1460,7 +1468,10 @@ private void replay(ApiMessage message, Optional snapshotId, lon */ private final RecordRedactor recordRedactor; - private final EventPerformanceMonitor eventPerformanceMonitor; + /** + * Monitors the performance of controller events and generates logs about it. + */ + private final EventPerformanceMonitor performanceMonitor; private QuorumController( FaultHandler nonFatalFaultHandler, @@ -1494,7 +1505,8 @@ private QuorumController( long delegationTokenExpiryCheckIntervalMs, long uncleanLeaderElectionCheckIntervalMs, String interBrokerListenerName, - int minSlowEventTimeMs + long controllerPerformanceSamplePeriodMs, + long controllerPerformanceAlwaysLogThresholdMs ) { this.nonFatalFaultHandler = nonFatalFaultHandler; this.fatalFaultHandler = fatalFaultHandler; @@ -1591,8 +1603,11 @@ private QuorumController( this.metaLogListener = new QuorumMetaLogListener(); this.curClaimEpoch = -1; this.recordRedactor = new RecordRedactor(configSchema); - this.eventPerformanceMonitor = new EventPerformanceMonitor(minSlowEventTimeMs, - controllerMetrics::getEventQueueProcessingTime99, logContext); + this.performanceMonitor = new EventPerformanceMonitor.Builder(). + setLogContext(logContext). + setPeriodNs(TimeUnit.MILLISECONDS.toNanos(controllerPerformanceSamplePeriodMs)). + setAlwaysLogThresholdNs(TimeUnit.MILLISECONDS.toNanos(controllerPerformanceAlwaysLogThresholdMs)). + build(); if (maxIdleIntervalNs.isPresent()) { registerWriteNoOpRecord(maxIdleIntervalNs.getAsLong()); } @@ -1606,7 +1621,7 @@ private QuorumController( } registerElectUnclean(TimeUnit.MILLISECONDS.toNanos(uncleanLeaderElectionCheckIntervalMs)); registerExpireDelegationTokens(MILLISECONDS.toNanos(delegationTokenExpiryCheckIntervalMs)); - registerUpdateSlowEventLogger(SECONDS.toNanos(30)); + registerGeneratePeriodicPerformanceMessage(); // OffsetControlManager must be initialized last, because its constructor will take the // initial in-memory snapshot of all extant timeline data structures. this.offsetControl = new OffsetControlManager.Builder(). @@ -1640,16 +1655,6 @@ private void registerWriteNoOpRecord(long maxIdleIntervalNs) { EnumSet.noneOf(PeriodicTaskFlag.class))); } - private void registerUpdateSlowEventLogger(long maxSlowEventWindowNs) { - periodicControl.registerTask(new PeriodicTask("updateSlowEventLoggerP99", - () -> { - eventPerformanceMonitor.refreshPercentile(); - return ControllerResult.of(Collections.emptyList(), false); - }, - maxSlowEventWindowNs, - EnumSet.noneOf(PeriodicTaskFlag.class))); - } - /** * Calculate what the period should be for the maybeFenceStaleBroker task. * @@ -1709,6 +1714,21 @@ private void registerElectUnclean(long checkIntervalNs) { EnumSet.of(PeriodicTaskFlag.VERBOSE))); } + /** + * Register the generatePeriodicPerformanceMessage task. + * + * This task periodically logs some statistics about controller performance. + */ + private void registerGeneratePeriodicPerformanceMessage() { + periodicControl.registerTask(new PeriodicTask("generatePeriodicPerformanceMessage", + () -> { + performanceMonitor.generatePeriodicPerformanceMessage(); + return ControllerResult.of(Collections.emptyList(), false); + }, + performanceMonitor.periodNs(), + EnumSet.noneOf(PeriodicTaskFlag.class))); + } + /** * Register the delegation token expiration task. * diff --git a/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java b/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java index e27bd3b83d9e7..10c3807da854b 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java +++ b/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java @@ -164,15 +164,6 @@ public void updateEventQueueProcessingTime(long durationMs) { eventQueueProcessingTimeUpdater.accept(durationMs); } - public double getEventQueueProcessingTime99() { - if (registry.isPresent()) { - Histogram histogram = registry.get().newHistogram(EVENT_QUEUE_PROCESSING_TIME_MS, true); - return histogram.getSnapshot().get99thPercentile(); - } else { - // Only returned in unit tests when a metrics registry is not set. - return 0.0; - } - } public void setLastAppliedRecordOffset(long offset) { lastAppliedRecordOffset.set(offset); } diff --git a/metadata/src/test/java/org/apache/kafka/controller/EventPerformanceMonitorTest.java b/metadata/src/test/java/org/apache/kafka/controller/EventPerformanceMonitorTest.java index 814f4f36a900b..b71a8e2fe92e4 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/EventPerformanceMonitorTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/EventPerformanceMonitorTest.java @@ -17,73 +17,104 @@ package org.apache.kafka.controller; -import org.apache.kafka.common.utils.LogContext; - import org.junit.jupiter.api.Test; -import java.util.concurrent.atomic.AtomicReference; +import java.util.AbstractMap; import static java.util.concurrent.TimeUnit.MILLISECONDS; -import static org.junit.jupiter.api.Assertions.assertFalse; -import static org.junit.jupiter.api.Assertions.assertTrue; - +import static java.util.concurrent.TimeUnit.SECONDS; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; public class EventPerformanceMonitorTest { @Test - public void testSlowEvents() { - LogContext logContext = new LogContext(); - - AtomicReference p99 = new AtomicReference<>(0.0); - EventPerformanceMonitor logger = new EventPerformanceMonitor(100, p99::get, logContext); - - // Initially, the p99 is zero - assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(10))); - assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(99))); - assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(100))); - - - // Idle controller, low p99 - p99.set(30.0); - logger.refreshPercentile(); - assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(90))); - assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(99))); - assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(100))); - - // Busy controller, high p99 - p99.set(1000.0); - logger.refreshPercentile(); - assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(100))); - assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(200))); - assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(1000))); - assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(2000))); + public void testDefaultPeriodNs() { + assertEquals(SECONDS.toNanos(60), + new EventPerformanceMonitor.Builder().build().periodNs()); + } + + @Test + public void testSlowestEventWithNoEvents() { + EventPerformanceMonitor monitor = new EventPerformanceMonitor.Builder().build(); + assertEquals(new AbstractMap.SimpleImmutableEntry<>(null, 0L), + monitor.slowestEvent()); + } + + @Test + public void testSlowestEventWithThreeEvents() { + EventPerformanceMonitor monitor = new EventPerformanceMonitor.Builder().build(); + monitor.observeEvent("fastEvent", MILLISECONDS.toNanos(2)); + monitor.observeEvent("slowEvent", MILLISECONDS.toNanos(100)); + assertEquals(new AbstractMap.SimpleImmutableEntry<>("slowEvent", MILLISECONDS.toNanos(100)), + monitor.slowestEvent()); + } + + @Test + public void testLogSlowEvent() { + EventPerformanceMonitor monitor = new EventPerformanceMonitor.Builder().build(); + assertEquals("Exceptionally slow controller event slowEvent took 5000 ms.", + monitor.doObserveEvent("slowEvent", SECONDS.toNanos(5))); + } + + @Test + public void testDoNotLogFastEvent() { + EventPerformanceMonitor monitor = new EventPerformanceMonitor.Builder().build(); + assertNull(monitor.doObserveEvent("slowEvent", MILLISECONDS.toNanos(250))); + } + + @Test + public void testNanosecondsToDecimalMillisWithZero() { + assertEquals("0.00", + EventPerformanceMonitor.nanosecondsToDecimalMillis(0)); + } + + @Test + public void testNanosecondsToDecimalMillisWith100() { + assertEquals("100.00", + EventPerformanceMonitor.nanosecondsToDecimalMillis(MILLISECONDS.toNanos(100))); + } + + @Test + public void testNanosecondsToDecimalMillisWith123456789() { + assertEquals("123.46", + EventPerformanceMonitor.nanosecondsToDecimalMillis(123456789)); + } + + @Test + public void testPeriodicPerformanceMessageWithNoEvents() { + EventPerformanceMonitor monitor = new EventPerformanceMonitor.Builder().build(); + assertEquals("In the last 60000 ms period, there were no controller events completed.", + monitor.periodicPerformanceMessage()); + } + + @Test + public void testPeriodicPerformanceMessageWithOneEvent() { + EventPerformanceMonitor monitor = new EventPerformanceMonitor.Builder().build(); + monitor.observeEvent("myEvent", MILLISECONDS.toNanos(12)); + assertEquals("In the last 60000 ms period, 1 controller events were completed, which took an " + + "average of 12.00 ms each. The slowest event was myEvent, which took 12.00 ms.", + monitor.periodicPerformanceMessage()); + } + + @Test + public void testPeriodicPerformanceMessageWithThreeEvents() { + EventPerformanceMonitor monitor = new EventPerformanceMonitor.Builder().build(); + monitor.observeEvent("myEvent", MILLISECONDS.toNanos(12)); + monitor.observeEvent("myEvent2", MILLISECONDS.toNanos(19)); + monitor.observeEvent("myEvent3", MILLISECONDS.toNanos(1)); + assertEquals("In the last 60000 ms period, 3 controller events were completed, which took an " + + "average of 10.67 ms each. The slowest event was myEvent2, which took 19.00 ms.", + monitor.periodicPerformanceMessage()); } @Test - public void testThresholdDisabled() { - LogContext logContext = new LogContext(); - - AtomicReference p99 = new AtomicReference<>(0.0); - // Set min slow event time to zero, effectively disabling the threshold - EventPerformanceMonitor logger = new EventPerformanceMonitor(0, p99::get, logContext); - - assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(0))); - assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(10))); - assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(99))); - assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(100))); - - p99.set(30.0); - logger.refreshPercentile(); - assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(0))); - assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(29))); - assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(30))); - assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(100))); - - - p99.set(1000.0); - logger.refreshPercentile(); - assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(100))); - assertFalse(logger.observeEvent("test", MILLISECONDS.toNanos(999))); - assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(1000))); - assertTrue(logger.observeEvent("test", MILLISECONDS.toNanos(2000))); + public void testGeneratePeriodicPerformanceMessageResetsState() { + EventPerformanceMonitor monitor = new EventPerformanceMonitor.Builder().build(); + monitor.observeEvent("myEvent", MILLISECONDS.toNanos(12)); + monitor.observeEvent("myEvent2", MILLISECONDS.toNanos(19)); + monitor.observeEvent("myEvent3", MILLISECONDS.toNanos(1)); + monitor.generatePeriodicPerformanceMessage(); + assertEquals("In the last 60000 ms period, there were no controller events completed.", + monitor.periodicPerformanceMessage()); } } diff --git a/server/src/main/java/org/apache/kafka/server/config/KRaftConfigs.java b/server/src/main/java/org/apache/kafka/server/config/KRaftConfigs.java index 6ad4647796416..a693471b7b3f5 100644 --- a/server/src/main/java/org/apache/kafka/server/config/KRaftConfigs.java +++ b/server/src/main/java/org/apache/kafka/server/config/KRaftConfigs.java @@ -114,9 +114,13 @@ public class KRaftConfigs { public static final String SERVER_MAX_STARTUP_TIME_MS_DOC = "The maximum number of milliseconds we will wait for the server to come up. " + "By default there is no limit. This should be used for testing only."; - public static final String MIN_SLOW_EVENT_TIME_MS_CONFIG = "controller.slow.event.min.ms"; - public static final int MIN_SLOW_EVENT_TIME_MS_DEFAULT = 200; - public static final String MIN_SLOW_EVENT_TIME_MS_DOC = "Log controller events with a p99 duration slower than this amount."; + public static final String CONTROLLER_PERFORMANCE_SAMPLE_PERIOD_MS = "controller.performance.sample.period.ms"; + public static final long CONTROLLER_PERFORMANCE_SAMPLE_PERIOD_MS_DEFAULT = 60000; + public static final String CONTROLLER_PERFORMANCE_SAMPLE_PERIOD_MS_DOC = "The number of milliseconds between periodic controller event performance log messages."; + + public static final String CONTROLLER_PERFORMANCE_ALWAYS_LOG_THRESHOLD_MS = "controller.performance.always.log.threshold.ms"; + public static final long CONTROLLER_PERFORMANCE_ALWAYS_LOG_THRESHOLD_MS_DEFAULT = 2000; + public static final String CONTROLLER_PERFORMANCE_ALWAYS_LOG_THRESHOLD_MS_DOC = "We will log an error message about controller events that take longer than this threshold."; public static final ConfigDef CONFIG_DEF = new ConfigDef() .define(METADATA_SNAPSHOT_MAX_NEW_RECORD_BYTES_CONFIG, LONG, METADATA_SNAPSHOT_MAX_NEW_RECORD_BYTES, atLeast(1), HIGH, METADATA_SNAPSHOT_MAX_NEW_RECORD_BYTES_DOC) @@ -135,6 +139,7 @@ public class KRaftConfigs { .define(METADATA_MAX_RETENTION_BYTES_CONFIG, LONG, METADATA_MAX_RETENTION_BYTES_DEFAULT, null, HIGH, METADATA_MAX_RETENTION_BYTES_DOC) .define(METADATA_MAX_RETENTION_MILLIS_CONFIG, LONG, LogConfig.DEFAULT_RETENTION_MS, null, HIGH, METADATA_MAX_RETENTION_MILLIS_DOC) .define(METADATA_MAX_IDLE_INTERVAL_MS_CONFIG, INT, METADATA_MAX_IDLE_INTERVAL_MS_DEFAULT, atLeast(0), LOW, METADATA_MAX_IDLE_INTERVAL_MS_DOC) - .defineInternal(MIN_SLOW_EVENT_TIME_MS_CONFIG, INT, MIN_SLOW_EVENT_TIME_MS_DEFAULT, atLeast(0), MEDIUM, MIN_SLOW_EVENT_TIME_MS_DOC) + .defineInternal(CONTROLLER_PERFORMANCE_SAMPLE_PERIOD_MS, LONG, CONTROLLER_PERFORMANCE_SAMPLE_PERIOD_MS_DEFAULT, atLeast(100), MEDIUM, CONTROLLER_PERFORMANCE_SAMPLE_PERIOD_MS_DOC) + .defineInternal(CONTROLLER_PERFORMANCE_ALWAYS_LOG_THRESHOLD_MS, LONG, CONTROLLER_PERFORMANCE_ALWAYS_LOG_THRESHOLD_MS_DEFAULT, atLeast(0), MEDIUM, CONTROLLER_PERFORMANCE_ALWAYS_LOG_THRESHOLD_MS_DOC) .defineInternal(SERVER_MAX_STARTUP_TIME_MS_CONFIG, LONG, SERVER_MAX_STARTUP_TIME_MS_DEFAULT, atLeast(0), MEDIUM, SERVER_MAX_STARTUP_TIME_MS_DOC); } From dee71c40dcd2c3f03c395b632f627ee7d6b3f1a8 Mon Sep 17 00:00:00 2001 From: "Colin P. McCabe" Date: Mon, 6 Jan 2025 12:53:02 -0800 Subject: [PATCH 9/9] nanosecondsToDecimalMillis -> formatNsAsDecimalMs --- .../kafka/controller/EventPerformanceMonitor.java | 6 +++--- .../controller/EventPerformanceMonitorTest.java | 12 ++++++------ 2 files changed, 9 insertions(+), 9 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/controller/EventPerformanceMonitor.java b/metadata/src/main/java/org/apache/kafka/controller/EventPerformanceMonitor.java index 5f0b2635122ee..fbe8b1c3cbbb8 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/EventPerformanceMonitor.java +++ b/metadata/src/main/java/org/apache/kafka/controller/EventPerformanceMonitor.java @@ -189,10 +189,10 @@ String periodicPerformanceMessage() { bld.append("there were no controller events completed."); } else { bld.append(numEvents).append(" controller events were completed, which took an average of "); - bld.append(nanosecondsToDecimalMillis(totalEventDurationNs / numEvents)); + bld.append(formatNsAsDecimalMs(totalEventDurationNs / numEvents)); bld.append(" ms each. The slowest event was ").append(slowestEventName); bld.append(", which took "); - bld.append(nanosecondsToDecimalMillis(slowestEventDurationNs)); + bld.append(formatNsAsDecimalMs(slowestEventDurationNs)); bld.append(" ms."); } return bld.toString(); @@ -204,7 +204,7 @@ String periodicPerformanceMessage() { * @param durationNs The duration in nanoseconds. * @return The decimal duration in milliseconds. */ - static String nanosecondsToDecimalMillis(long durationNs) { + static String formatNsAsDecimalMs(long durationNs) { double number = NANOSECONDS.toMicros(durationNs); number /= 1000; return MILLISECOND_DECIMAL_FORMAT.format(number); diff --git a/metadata/src/test/java/org/apache/kafka/controller/EventPerformanceMonitorTest.java b/metadata/src/test/java/org/apache/kafka/controller/EventPerformanceMonitorTest.java index b71a8e2fe92e4..81e01679dc0fd 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/EventPerformanceMonitorTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/EventPerformanceMonitorTest.java @@ -63,21 +63,21 @@ public void testDoNotLogFastEvent() { } @Test - public void testNanosecondsToDecimalMillisWithZero() { + public void testFormatNsAsDecimalMsWithZero() { assertEquals("0.00", - EventPerformanceMonitor.nanosecondsToDecimalMillis(0)); + EventPerformanceMonitor.formatNsAsDecimalMs(0)); } @Test - public void testNanosecondsToDecimalMillisWith100() { + public void testFormatNsAsDecimalMsWith100() { assertEquals("100.00", - EventPerformanceMonitor.nanosecondsToDecimalMillis(MILLISECONDS.toNanos(100))); + EventPerformanceMonitor.formatNsAsDecimalMs(MILLISECONDS.toNanos(100))); } @Test - public void testNanosecondsToDecimalMillisWith123456789() { + public void testFormatNsAsDecimalMsWith123456789() { assertEquals("123.46", - EventPerformanceMonitor.nanosecondsToDecimalMillis(123456789)); + EventPerformanceMonitor.formatNsAsDecimalMs(123456789)); } @Test