diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java index 82e1f8b93da3c..bb75a40657311 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java @@ -20,7 +20,7 @@ import org.apache.kafka.clients.consumer.internals.ConsumerDelegate; import org.apache.kafka.clients.consumer.internals.ConsumerDelegateCreator; import org.apache.kafka.clients.consumer.internals.ConsumerMetadata; -import org.apache.kafka.clients.consumer.internals.KafkaConsumerMetrics; +import org.apache.kafka.clients.consumer.internals.metrics.KafkaConsumerMetrics; import org.apache.kafka.clients.consumer.internals.SubscriptionState; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.Metric; diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java index b61202ea7b609..31481079cc9e4 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java @@ -60,7 +60,8 @@ import org.apache.kafka.clients.consumer.internals.events.TopicMetadataApplicationEvent; import org.apache.kafka.clients.consumer.internals.events.UnsubscribeApplicationEvent; import org.apache.kafka.clients.consumer.internals.events.ValidatePositionsApplicationEvent; -import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetrics; +import org.apache.kafka.clients.consumer.internals.metrics.KafkaConsumerMetrics; +import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetricsManager; import org.apache.kafka.common.Cluster; import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.KafkaException; @@ -396,7 +397,7 @@ private void process(final ConsumerRebalanceListenerCallbackNeededEvent event) { logContext, subscriptions, time, - new RebalanceCallbackMetrics(metrics) + new RebalanceCallbackMetricsManager(metrics) ); this.backgroundEventProcessor = new BackgroundEventProcessor( logContext, @@ -540,7 +541,7 @@ private void process(final ConsumerRebalanceListenerCallbackNeededEvent event) { logContext, subscriptions, time, - new RebalanceCallbackMetrics(metrics) + new RebalanceCallbackMetricsManager(metrics) ); ApiVersions apiVersions = new ApiVersions(); Supplier networkClientDelegateSupplier = () -> new NetworkClientDelegate( diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/CommitRequestManager.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/CommitRequestManager.java index 75a040d293335..2ed6d6ca7da17 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/CommitRequestManager.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/CommitRequestManager.java @@ -21,6 +21,7 @@ import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.clients.consumer.RetriableCommitFailedException; +import org.apache.kafka.clients.consumer.internals.metrics.OffsetCommitMetricsManager; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.DisconnectException; @@ -32,11 +33,6 @@ import org.apache.kafka.common.message.OffsetCommitRequestData; import org.apache.kafka.common.message.OffsetCommitResponseData; import org.apache.kafka.common.metrics.Metrics; -import org.apache.kafka.common.metrics.Sensor; -import org.apache.kafka.common.metrics.stats.Avg; -import org.apache.kafka.common.metrics.stats.Max; -import org.apache.kafka.common.metrics.stats.Meter; -import org.apache.kafka.common.metrics.stats.WindowedCount; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.record.RecordBatch; import org.apache.kafka.common.requests.AbstractRequest; @@ -66,8 +62,6 @@ import java.util.function.BiConsumer; import java.util.stream.Collectors; -import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.CONSUMER_METRIC_GROUP_PREFIX; -import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.COORDINATOR_METRICS_SUFFIX; import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.THROW_ON_FETCH_STABLE_OFFSET_UNSUPPORTED; import static org.apache.kafka.clients.consumer.internals.NetworkClientDelegate.PollResult.EMPTY; import static org.apache.kafka.common.protocol.Errors.COORDINATOR_LOAD_IN_PROGRESS; @@ -79,6 +73,7 @@ public class CommitRequestManager implements RequestManager, MemberStateListener private final Optional autoCommitState; private final CoordinatorRequestManager coordinatorRequestManager; private final OffsetCommitCallbackInvoker offsetCommitCallbackInvoker; + private final OffsetCommitMetricsManager metricsManager; private final long retryBackoffMs; private final String groupId; private final Optional groupInstanceId; @@ -88,7 +83,6 @@ public class CommitRequestManager implements RequestManager, MemberStateListener private final boolean throwOnFetchStableOffsetUnsupported; final PendingRequests pendingRequests; private boolean closing = false; - private Sensor commitSensor; /** * Latest member ID and epoch received via the {@link #onMemberEpochUpdated(Optional, Optional)}, @@ -119,7 +113,6 @@ public CommitRequestManager( config.getLong(ConsumerConfig.RETRY_BACKOFF_MS_CONFIG), config.getLong(ConsumerConfig.RETRY_BACKOFF_MAX_MS_CONFIG), OptionalDouble.empty(), - CONSUMER_METRIC_GROUP_PREFIX, metrics); } @@ -136,7 +129,6 @@ public CommitRequestManager( final long retryBackoffMs, final long retryBackoffMaxMs, final OptionalDouble jitter, - final String metricGroupPrefix, final Metrics metrics) { Objects.requireNonNull(coordinatorRequestManager, "Coordinator is needed upon committing offsets"); this.logContext = logContext; @@ -158,8 +150,8 @@ public CommitRequestManager( this.jitter = jitter; this.throwOnFetchStableOffsetUnsupported = config.getBoolean(THROW_ON_FETCH_STABLE_OFFSET_UNSUPPORTED); this.memberInfo = new MemberInfo(); + this.metricsManager = new OffsetCommitMetricsManager(metrics); this.offsetCommitCallbackInvoker = offsetCommitCallbackInvoker; - this.commitSensor = addCommitSensor(metrics, metricGroupPrefix); } /** @@ -400,23 +392,6 @@ public NetworkClientDelegate.PollResult drainPendingOffsetCommitRequests() { return new NetworkClientDelegate.PollResult(Long.MAX_VALUE, requests); } - private Sensor addCommitSensor(Metrics metrics, String metricGrpPrefix) { - String metricGrpName = metricGrpPrefix + COORDINATOR_METRICS_SUFFIX; - Sensor sensor = metrics.sensor("commit-latency"); - sensor.add(metrics.metricName("commit-latency-avg", - metricGrpName, - "The average time taken for a commit request"), new Avg()); - sensor.add(metrics.metricName("commit-latency-max", - metricGrpName, - "The max time taken for a commit request"), new Max()); - sensor.add(new Meter(new WindowedCount(), - metrics.metricName("commit-rate", metricGrpName, - "The number of commit calls per second"), - metrics.metricName("commit-total", metricGrpName, - "The total number of commit calls"))); - return sensor; - } - private class OffsetCommitRequestState extends RetriableRequestState { private final Map offsets; private final String groupId; @@ -511,8 +486,8 @@ public NetworkClientDelegate.UnsentRequest toUnsentRequest() { * - fail the future with a non-recoverable KafkaException for all unexpected errors (even if retriable) */ @Override - void onResponse(final ClientResponse response) { - commitSensor.record(response.requestLatencyMs()); + public void onResponse(final ClientResponse response) { + metricsManager.recordRequestLatency(response.requestLatencyMs()); long currentTimeMs = response.receivedTimeMs(); OffsetCommitResponse commitResponse = (OffsetCommitResponse) response.responseBody(); Set unauthorizedTopics = new HashSet<>(); diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java index 398baf7fb079c..5af8633806da8 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java @@ -29,7 +29,7 @@ import org.apache.kafka.clients.consumer.OffsetCommitCallback; import org.apache.kafka.clients.consumer.RetriableCommitFailedException; import org.apache.kafka.clients.consumer.internals.Utils.TopicPartitionComparator; -import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetrics; +import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetricsManager; import org.apache.kafka.common.Cluster; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.Node; @@ -235,7 +235,7 @@ public ConsumerCoordinator(GroupRebalanceConfig rebalanceConfig, logContext, subscriptions, time, - new RebalanceCallbackMetrics(metrics, metricGrpPrefix) + new RebalanceCallbackMetricsManager(metrics, metricGrpPrefix) ); this.metadata.requestUpdate(true); } diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerDelegate.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerDelegate.java index 612827ebe83df..67e8d991d639f 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerDelegate.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerDelegate.java @@ -17,6 +17,7 @@ package org.apache.kafka.clients.consumer.internals; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.internals.metrics.KafkaConsumerMetrics; import org.apache.kafka.common.metrics.Metrics; import org.apache.kafka.common.utils.Timer; diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerRebalanceListenerInvoker.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerRebalanceListenerInvoker.java index ea60436deb1ed..dcdd303fe95d1 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerRebalanceListenerInvoker.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerRebalanceListenerInvoker.java @@ -17,7 +17,7 @@ package org.apache.kafka.clients.consumer.internals; import org.apache.kafka.clients.consumer.ConsumerRebalanceListener; -import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetrics; +import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetricsManager; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.InterruptException; import org.apache.kafka.common.errors.WakeupException; @@ -41,16 +41,16 @@ public class ConsumerRebalanceListenerInvoker { private final Logger log; private final SubscriptionState subscriptions; private final Time time; - private final RebalanceCallbackMetrics metrics; + private final RebalanceCallbackMetricsManager metricsManager; ConsumerRebalanceListenerInvoker(LogContext logContext, SubscriptionState subscriptions, Time time, - RebalanceCallbackMetrics metrics) { + RebalanceCallbackMetricsManager metricsManager) { this.log = logContext.logger(getClass()); this.subscriptions = subscriptions; this.time = time; - this.metrics = metrics; + this.metricsManager = metricsManager; } public Exception invokePartitionsAssigned(final SortedSet assignedPartitions) { @@ -62,7 +62,7 @@ public Exception invokePartitionsAssigned(final SortedSet assign try { final long startMs = time.milliseconds(); listener.get().onPartitionsAssigned(assignedPartitions); - metrics.assignCallbackSensor.record(time.milliseconds() - startMs); + metricsManager.recordPartitionsAssignedLatency(time.milliseconds() - startMs); } catch (WakeupException | InterruptException e) { throw e; } catch (Exception e) { @@ -88,7 +88,7 @@ public Exception invokePartitionsRevoked(final SortedSet revoked try { final long startMs = time.milliseconds(); listener.get().onPartitionsRevoked(revokedPartitions); - metrics.revokeCallbackSensor.record(time.milliseconds() - startMs); + metricsManager.recordPartitionsRevokedLatency(time.milliseconds() - startMs); } catch (WakeupException | InterruptException e) { throw e; } catch (Exception e) { @@ -114,7 +114,7 @@ public Exception invokePartitionsLost(final SortedSet lostPartit try { final long startMs = time.milliseconds(); listener.get().onPartitionsLost(lostPartitions); - metrics.loseCallbackSensor.record(time.milliseconds() - startMs); + metricsManager.recordPartitionsLostLatency(time.milliseconds() - startMs); } catch (WakeupException | InterruptException e) { throw e; } catch (Exception e) { diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerUtils.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerUtils.java index bc7590c38f819..5a7745b40dd0d 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerUtils.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerUtils.java @@ -66,6 +66,7 @@ public final class ConsumerUtils { public static final String CONSUMER_JMX_PREFIX = "kafka.consumer"; public static final String CONSUMER_METRIC_GROUP_PREFIX = "consumer"; public static final String COORDINATOR_METRICS_SUFFIX = "-coordinator-metrics"; + public static final String CONSUMER_METRICS_SUFFIX = "-metrics"; /** * A fixed, large enough value will suffice for max. diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/HeartbeatRequestManager.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/HeartbeatRequestManager.java index 11846eedeafdc..03e11ddfa02b4 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/HeartbeatRequestManager.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/HeartbeatRequestManager.java @@ -23,10 +23,12 @@ import org.apache.kafka.clients.consumer.internals.events.BackgroundEventHandler; import org.apache.kafka.clients.consumer.internals.events.ErrorBackgroundEvent; import org.apache.kafka.clients.consumer.internals.events.GroupMetadataUpdateEvent; +import org.apache.kafka.clients.consumer.internals.metrics.HeartbeatMetricsManager; import org.apache.kafka.common.Uuid; import org.apache.kafka.common.errors.GroupAuthorizationException; import org.apache.kafka.common.errors.RetriableException; import org.apache.kafka.common.message.ConsumerGroupHeartbeatRequestData; +import org.apache.kafka.common.metrics.Metrics; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.requests.ConsumerGroupHeartbeatRequest; import org.apache.kafka.common.requests.ConsumerGroupHeartbeatResponse; @@ -108,9 +110,13 @@ public class HeartbeatRequestManager implements RequestManager { * sending heartbeat until the next poll. */ private final Timer pollTimer; - private GroupMetadataUpdateEvent previousGroupMetadataUpdateEvent = null; + /** + * Holding the heartbeat sensor to measure heartbeat timing and response latency + */ + private final HeartbeatMetricsManager metricsManager; + public HeartbeatRequestManager( final LogContext logContext, final Time time, @@ -118,7 +124,8 @@ public HeartbeatRequestManager( final CoordinatorRequestManager coordinatorRequestManager, final SubscriptionState subscriptions, final MembershipManager membershipManager, - final BackgroundEventHandler backgroundEventHandler) { + final BackgroundEventHandler backgroundEventHandler, + final Metrics metrics) { this.coordinatorRequestManager = coordinatorRequestManager; this.logger = logContext.logger(getClass()); this.membershipManager = membershipManager; @@ -130,6 +137,7 @@ public HeartbeatRequestManager( this.heartbeatRequestState = new HeartbeatRequestState(logContext, time, 0, retryBackoffMs, retryBackoffMaxMs, maxPollIntervalMs); this.pollTimer = time.timer(maxPollIntervalMs); + this.metricsManager = new HeartbeatMetricsManager(metrics); } // Visible for testing @@ -141,7 +149,8 @@ public HeartbeatRequestManager( final MembershipManager membershipManager, final HeartbeatState heartbeatState, final HeartbeatRequestState heartbeatRequestState, - final BackgroundEventHandler backgroundEventHandler) { + final BackgroundEventHandler backgroundEventHandler, + final Metrics metrics) { this.logger = logContext.logger(this.getClass()); this.maxPollIntervalMs = config.getInt(CommonClientConfigs.MAX_POLL_INTERVAL_MS_CONFIG); this.coordinatorRequestManager = coordinatorRequestManager; @@ -150,6 +159,7 @@ public HeartbeatRequestManager( this.membershipManager = membershipManager; this.backgroundEventHandler = backgroundEventHandler; this.pollTimer = timer; + this.metricsManager = new HeartbeatMetricsManager(metrics); } /** @@ -245,6 +255,7 @@ private NetworkClientDelegate.UnsentRequest makeHeartbeatRequest(final long curr NetworkClientDelegate.UnsentRequest request = makeHeartbeatRequest(ignoreResponse); heartbeatRequestState.onSendAttempt(currentTimeMs); membershipManager.onHeartbeatRequestSent(); + metricsManager.recordHeartbeatSentMs(currentTimeMs); return request; } @@ -257,6 +268,7 @@ private NetworkClientDelegate.UnsentRequest makeHeartbeatRequest(final boolean i else return request.whenComplete((response, exception) -> { if (response != null) { + metricsManager.recordRequestLatency(response.requestLatencyMs()); onResponse((ConsumerGroupHeartbeatResponse) response.responseBody(), request.handler().completionTimeMs()); } else { onFailure(exception, request.handler().completionTimeMs()); @@ -267,6 +279,7 @@ private NetworkClientDelegate.UnsentRequest makeHeartbeatRequest(final boolean i private NetworkClientDelegate.UnsentRequest logResponse(final NetworkClientDelegate.UnsentRequest request) { return request.whenComplete((response, exception) -> { if (response != null) { + metricsManager.recordRequestLatency(response.requestLatencyMs()); Errors error = Errors.forCode(((ConsumerGroupHeartbeatResponse) response.responseBody()).data().errorCode()); if (error == Errors.NONE) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/LegacyKafkaConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/LegacyKafkaConsumer.java index bcc8cb40f2d91..80be0959d69cb 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/LegacyKafkaConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/LegacyKafkaConsumer.java @@ -35,6 +35,7 @@ import org.apache.kafka.clients.consumer.OffsetAndTimestamp; import org.apache.kafka.clients.consumer.OffsetCommitCallback; import org.apache.kafka.clients.consumer.OffsetResetStrategy; +import org.apache.kafka.clients.consumer.internals.metrics.KafkaConsumerMetrics; import org.apache.kafka.common.Cluster; import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.KafkaException; diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/RequestManagers.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/RequestManagers.java index 6ca2e398de382..b3c0c5b63b58e 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/RequestManagers.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/RequestManagers.java @@ -194,7 +194,8 @@ protected RequestManagers create() { coordinator, subscriptions, membershipManager, - backgroundEventHandler); + backgroundEventHandler, + metrics); } return new RequestManagers( diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/AbstractConsumerMetrics.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/AbstractConsumerMetrics.java deleted file mode 100644 index 2273bfd494a3b..0000000000000 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/AbstractConsumerMetrics.java +++ /dev/null @@ -1,29 +0,0 @@ -/* - * 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.clients.consumer.internals.metrics; - -import java.util.Optional; - -import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.CONSUMER_METRIC_GROUP_PREFIX; -import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.COORDINATOR_METRICS_SUFFIX; - -public abstract class AbstractConsumerMetrics { - protected String groupMetricsName = CONSUMER_METRIC_GROUP_PREFIX + COORDINATOR_METRICS_SUFFIX; - public AbstractConsumerMetrics(Optional grpMetricsPrefix) { - grpMetricsPrefix.ifPresent(s -> this.groupMetricsName = s + COORDINATOR_METRICS_SUFFIX); - } -} diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/HeartbeatMetricsManager.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/HeartbeatMetricsManager.java new file mode 100644 index 0000000000000..a5654dfc85532 --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/HeartbeatMetricsManager.java @@ -0,0 +1,77 @@ +/* + * 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.clients.consumer.internals.metrics; + +import org.apache.kafka.common.MetricName; +import org.apache.kafka.common.metrics.Measurable; +import org.apache.kafka.common.metrics.Metrics; +import org.apache.kafka.common.metrics.Sensor; +import org.apache.kafka.common.metrics.stats.Max; +import org.apache.kafka.common.metrics.stats.Meter; +import org.apache.kafka.common.metrics.stats.WindowedCount; + +import java.util.concurrent.TimeUnit; + +import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.CONSUMER_METRIC_GROUP_PREFIX; +import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.COORDINATOR_METRICS_SUFFIX; + +public class HeartbeatMetricsManager { + // MetricName visible for testing + final MetricName heartbeatResponseTimeMax; + final MetricName heartbeatRate; + final MetricName heartbeatTotal; + final MetricName lastHeartbeatSecondsAgo; + private final Sensor heartbeatSensor; + private long lastHeartbeatMs = -1L; + + public HeartbeatMetricsManager(Metrics metrics) { + final String metricGroupName = CONSUMER_METRIC_GROUP_PREFIX + COORDINATOR_METRICS_SUFFIX; + heartbeatSensor = metrics.sensor("heartbeat-latency"); + heartbeatResponseTimeMax = metrics.metricName("heartbeat-response-time-max", + metricGroupName, + "The max time taken to receive a response to a heartbeat request"); + heartbeatSensor.add(heartbeatResponseTimeMax, new Max()); + + // windowed meters + heartbeatRate = metrics.metricName("heartbeat-rate", metricGroupName, "The number of heartbeats per second"); + heartbeatTotal = metrics.metricName("heartbeat-total", metricGroupName, "The total number of heartbeats"); + heartbeatSensor.add(new Meter(new WindowedCount(), + heartbeatRate, + heartbeatTotal)); + + Measurable lastHeartbeat = (config, now) -> { + final long lastHeartbeatSend = lastHeartbeatMs; + if (lastHeartbeatSend < 0L) + // if no heartbeat is ever triggered, just return -1. + return -1d; + else + return TimeUnit.SECONDS.convert(now - lastHeartbeatSend, TimeUnit.MILLISECONDS); + }; + lastHeartbeatSecondsAgo = metrics.metricName("last-heartbeat-seconds-ago", + metricGroupName, + "The number of seconds since the last coordinator heartbeat was sent"); + metrics.addMetric(lastHeartbeatSecondsAgo, lastHeartbeat); + } + + public void recordHeartbeatSentMs(long timeMs) { + lastHeartbeatMs = timeMs; + } + + public void recordRequestLatency(long requestLatencyMs) { + heartbeatSensor.record(requestLatencyMs); + } +} diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/KafkaConsumerMetrics.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/KafkaConsumerMetrics.java similarity index 93% rename from clients/src/main/java/org/apache/kafka/clients/consumer/internals/KafkaConsumerMetrics.java rename to clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/KafkaConsumerMetrics.java index 0dc8a33b11bbb..52502e714a947 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/KafkaConsumerMetrics.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/KafkaConsumerMetrics.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.clients.consumer.internals; +package org.apache.kafka.clients.consumer.internals.metrics; import org.apache.kafka.common.MetricName; import org.apache.kafka.common.metrics.Measurable; @@ -26,20 +26,22 @@ import java.util.concurrent.TimeUnit; +import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.CONSUMER_METRICS_SUFFIX; + public class KafkaConsumerMetrics implements AutoCloseable { + private final Metrics metrics; private final MetricName lastPollMetricName; private final Sensor timeBetweenPollSensor; private final Sensor pollIdleSensor; private final Sensor committedSensor; private final Sensor commitSyncSensor; - private final Metrics metrics; private long lastPollMs; private long pollStartMs; private long timeSinceLastPollMs; public KafkaConsumerMetrics(Metrics metrics, String metricGrpPrefix) { this.metrics = metrics; - String metricGroupName = metricGrpPrefix + "-metrics"; + final String metricGroupName = metricGrpPrefix + CONSUMER_METRICS_SUFFIX; Measurable lastPoll = (mConfig, now) -> { if (lastPollMs == 0L) // if no poll is ever triggered, just return -1. @@ -48,7 +50,7 @@ public KafkaConsumerMetrics(Metrics metrics, String metricGrpPrefix) { return TimeUnit.SECONDS.convert(now - lastPollMs, TimeUnit.MILLISECONDS); }; this.lastPollMetricName = metrics.metricName("last-poll-seconds-ago", - metricGroupName, "The number of seconds since the last poll() invocation."); + metricGroupName, "The number of seconds since the last poll() invocation."); metrics.addMetric(lastPollMetricName, lastPoll); this.timeBetweenPollSensor = metrics.sensor("time-between-poll"); diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/OffsetCommitMetricsManager.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/OffsetCommitMetricsManager.java new file mode 100644 index 0000000000000..d700299ef1801 --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/OffsetCommitMetricsManager.java @@ -0,0 +1,62 @@ +/* + * 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.clients.consumer.internals.metrics; + +import org.apache.kafka.common.MetricName; +import org.apache.kafka.common.metrics.Metrics; +import org.apache.kafka.common.metrics.Sensor; +import org.apache.kafka.common.metrics.stats.Avg; +import org.apache.kafka.common.metrics.stats.Max; +import org.apache.kafka.common.metrics.stats.Meter; +import org.apache.kafka.common.metrics.stats.WindowedCount; + +import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.CONSUMER_METRIC_GROUP_PREFIX; +import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.COORDINATOR_METRICS_SUFFIX; + +public class OffsetCommitMetricsManager { + final MetricName commitLatencyAvg; + final MetricName commitLatencyMax; + final MetricName commitRate; + final MetricName commitTotal; + private final Sensor commitSensor; + + public OffsetCommitMetricsManager(Metrics metrics) { + final String metricGroupName = CONSUMER_METRIC_GROUP_PREFIX + COORDINATOR_METRICS_SUFFIX; + commitSensor = metrics.sensor("commit-latency"); + commitLatencyAvg = metrics.metricName("commit-latency-avg", + metricGroupName, + "The average time taken for a commit request"); + commitSensor.add(commitLatencyAvg, new Avg()); + commitLatencyMax = metrics.metricName("commit-latency-max", + metricGroupName, + "The max time taken for a commit request"); + commitSensor.add(commitLatencyMax, new Max()); + commitRate = metrics.metricName("commit-rate", + metricGroupName, + "The number of commit calls per second"); + commitTotal = metrics.metricName("commit-total", + metricGroupName, + "The total number of commit calls"); + commitSensor.add(new Meter(new WindowedCount(), + commitRate, + commitTotal)); + } + + public void recordRequestLatency(long responseLatencyMs) { + this.commitSensor.record(responseLatencyMs); + } +} diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/RebalanceCallbackMetrics.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/RebalanceCallbackMetrics.java deleted file mode 100644 index a27581505c517..0000000000000 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/RebalanceCallbackMetrics.java +++ /dev/null @@ -1,61 +0,0 @@ -/* - * 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.clients.consumer.internals.metrics; - -import org.apache.kafka.common.metrics.Metrics; -import org.apache.kafka.common.metrics.Sensor; -import org.apache.kafka.common.metrics.stats.Avg; -import org.apache.kafka.common.metrics.stats.Max; - -import java.util.Optional; - -public class RebalanceCallbackMetrics extends AbstractConsumerMetrics { - public final Sensor revokeCallbackSensor; - public final Sensor assignCallbackSensor; - public final Sensor loseCallbackSensor; - - public RebalanceCallbackMetrics(Metrics metrics) { - this(metrics, null); - } - - public RebalanceCallbackMetrics(Metrics metrics, String grpMetricsPrefix) { - super(Optional.ofNullable(grpMetricsPrefix)); - revokeCallbackSensor = metrics.sensor("partition-revoked-latency"); - revokeCallbackSensor.add(metrics.metricName("partition-revoked-latency-avg", - groupMetricsName, - "The average time taken for a partition-revoked rebalance listener callback"), new Avg()); - revokeCallbackSensor.add(metrics.metricName("partition-revoked-latency-max", - groupMetricsName, - "The max time taken for a partition-revoked rebalance listener callback"), new Max()); - - assignCallbackSensor = metrics.sensor("partition-assigned-latency"); - assignCallbackSensor.add(metrics.metricName("partition-assigned-latency-avg", - groupMetricsName, - "The average time taken for a partition-assigned rebalance listener callback"), new Avg()); - assignCallbackSensor.add(metrics.metricName("partition-assigned-latency-max", - groupMetricsName, - "The max time taken for a partition-assigned rebalance listener callback"), new Max()); - - loseCallbackSensor = metrics.sensor("partition-lost-latency"); - loseCallbackSensor.add(metrics.metricName("partition-lost-latency-avg", - groupMetricsName, - "The average time taken for a partition-lost rebalance listener callback"), new Avg()); - loseCallbackSensor.add(metrics.metricName("partition-lost-latency-max", - groupMetricsName, - "The max time taken for a partition-lost rebalance listener callback"), new Max()); - } -} diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/RebalanceCallbackMetricsManager.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/RebalanceCallbackMetricsManager.java new file mode 100644 index 0000000000000..f70b891864f1c --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/RebalanceCallbackMetricsManager.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.clients.consumer.internals.metrics; + +import org.apache.kafka.common.MetricName; +import org.apache.kafka.common.metrics.Metrics; +import org.apache.kafka.common.metrics.Sensor; +import org.apache.kafka.common.metrics.stats.Avg; +import org.apache.kafka.common.metrics.stats.Max; + +import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.CONSUMER_METRIC_GROUP_PREFIX; +import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.COORDINATOR_METRICS_SUFFIX; + +public class RebalanceCallbackMetricsManager { + final MetricName partitionRevokeLatencyAvg; + final MetricName partitionAssignLatencyAvg; + final MetricName partitionLostLatencyAvg; + final MetricName partitionRevokeLatencyMax; + final MetricName partitionAssignLatencyMax; + final MetricName partitionLostLatencyMax; + private final Sensor partitionRevokeCallbackSensor; + private final Sensor partitionAssignCallbackSensor; + private final Sensor partitionLostCallbackSensor; + + public RebalanceCallbackMetricsManager(Metrics metrics) { + this(metrics, CONSUMER_METRIC_GROUP_PREFIX); + } + + public RebalanceCallbackMetricsManager(Metrics metrics, String grpMetricsPrefix) { + final String metricGroupName = grpMetricsPrefix + COORDINATOR_METRICS_SUFFIX; + partitionRevokeCallbackSensor = metrics.sensor("partition-revoked-latency"); + partitionRevokeLatencyAvg = metrics.metricName("partition-revoked-latency-avg", + metricGroupName, + "The average time taken for a partition-revoked rebalance listener callback"); + partitionRevokeCallbackSensor.add(partitionRevokeLatencyAvg, new Avg()); + partitionRevokeLatencyMax = metrics.metricName("partition-revoked-latency-max", + metricGroupName, + "The max time taken for a partition-revoked rebalance listener callback"); + partitionRevokeCallbackSensor.add(partitionRevokeLatencyMax, new Max()); + + partitionAssignCallbackSensor = metrics.sensor("partition-assigned-latency"); + partitionAssignLatencyAvg = metrics.metricName("partition-assigned-latency-avg", + metricGroupName, + "The average time taken for a partition-assigned rebalance listener callback"); + partitionAssignCallbackSensor.add(partitionAssignLatencyAvg, new Avg()); + partitionAssignLatencyMax = metrics.metricName("partition-assigned-latency-max", + metricGroupName, + "The max time taken for a partition-assigned rebalance listener callback"); + partitionAssignCallbackSensor.add(partitionAssignLatencyMax, new Max()); + + partitionLostCallbackSensor = metrics.sensor("partition-lost-latency"); + partitionLostLatencyAvg = metrics.metricName("partition-lost-latency-avg", + metricGroupName, + "The average time taken for a partition-lost rebalance listener callback"); + partitionLostCallbackSensor.add(partitionLostLatencyAvg, new Avg()); + partitionLostLatencyMax = metrics.metricName("partition-lost-latency-max", + metricGroupName, + "The max time taken for a partition-lost rebalance listener callback"); + partitionLostCallbackSensor.add(partitionLostLatencyMax, new Max()); + } + + public void recordPartitionsRevokedLatency(long latencyMs) { + partitionRevokeCallbackSensor.record(latencyMs); + } + + public void recordPartitionsAssignedLatency(long latencyMs) { + partitionAssignCallbackSensor.record(latencyMs); + } + + public void recordPartitionsLostLatency(long latencyMs) { + partitionLostCallbackSensor.record(latencyMs); + } +} diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/CommitRequestManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/CommitRequestManagerTest.java index b408960629796..15e1e0b4fb4bc 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/CommitRequestManagerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/CommitRequestManagerTest.java @@ -75,7 +75,6 @@ import static org.apache.kafka.clients.consumer.internals.ConsumerTestBuilder.DEFAULT_GROUP_ID; import static org.apache.kafka.clients.consumer.internals.ConsumerTestBuilder.DEFAULT_GROUP_INSTANCE_ID; import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.CONSUMER_METRIC_GROUP_PREFIX; -import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.COORDINATOR_METRICS_SUFFIX; import static org.apache.kafka.test.TestUtils.assertFutureThrows; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; @@ -94,7 +93,7 @@ public class CommitRequestManagerTest { private long retryBackoffMs = 100; private long retryBackoffMaxMs = 1000; private String consumerMetricGroupPrefix = CONSUMER_METRIC_GROUP_PREFIX; - private String consumerMetricGroupName = consumerMetricGroupPrefix + COORDINATOR_METRICS_SUFFIX; + private static final String CONSUMER_COORDINATOR_METRICS = "consumer-coordinator-metrics"; private Node mockedNode = new Node(1, "host1", 9092); private SubscriptionState subscriptionState; private LogContext logContext; @@ -159,6 +158,9 @@ public void testPollEnsureAutocommitSent() { 1, (short) 1, Errors.NONE))); + + assertEquals(0.03, (double) getMetric("commit-rate").metricValue(), 0.01); + assertEquals(1.0, getMetric("commit-total").metricValue()); } @Test @@ -1121,7 +1123,6 @@ private CommitRequestManager create(final boolean autoCommitEnabled, final long retryBackoffMs, retryBackoffMaxMs, OptionalDouble.of(0), - consumerMetricGroupPrefix, metrics)); } @@ -1246,6 +1247,6 @@ private ClientResponse buildOffsetFetchClientResponse( private KafkaMetric getMetric(String name) { return metrics.metrics().get(metrics.metricName( name, - consumerMetricGroupName)); + CONSUMER_COORDINATOR_METRICS)); } } diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerTestBuilder.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerTestBuilder.java index 768b150edebd4..f6f71f3f7a95b 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerTestBuilder.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerTestBuilder.java @@ -25,7 +25,7 @@ import org.apache.kafka.clients.consumer.internals.events.ApplicationEventProcessor; import org.apache.kafka.clients.consumer.internals.events.BackgroundEvent; import org.apache.kafka.clients.consumer.internals.events.BackgroundEventHandler; -import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetrics; +import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetricsManager; import org.apache.kafka.common.internals.ClusterResourceListeners; import org.apache.kafka.common.metrics.Metrics; import org.apache.kafka.common.requests.MetadataResponse; @@ -229,7 +229,8 @@ public ConsumerTestBuilder(Optional groupInfo, boolean enableA mm, heartbeatState, heartbeatRequestState, - backgroundEventHandler)); + backgroundEventHandler, + metrics)); this.coordinatorRequestManager = Optional.of(coordinator); this.commitRequestManager = Optional.of(commit); @@ -275,7 +276,7 @@ public ConsumerTestBuilder(Optional groupInfo, boolean enableA logContext, subscriptions, time, - new RebalanceCallbackMetrics(metrics) + new RebalanceCallbackMetricsManager(metrics) ); } diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/HeartbeatRequestManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/HeartbeatRequestManagerTest.java index e44414b0677b9..0df6fa94c39cb 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/HeartbeatRequestManagerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/HeartbeatRequestManagerTest.java @@ -28,6 +28,8 @@ import org.apache.kafka.common.errors.TimeoutException; import org.apache.kafka.common.message.ConsumerGroupHeartbeatRequestData; import org.apache.kafka.common.message.ConsumerGroupHeartbeatResponseData; +import org.apache.kafka.common.metrics.KafkaMetric; +import org.apache.kafka.common.metrics.Metrics; import org.apache.kafka.common.protocol.ApiKeys; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.requests.ConsumerGroupHeartbeatRequest; @@ -35,6 +37,7 @@ import org.apache.kafka.common.requests.RequestHeader; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.utils.LogContext; +import org.apache.kafka.common.utils.MockTime; import org.apache.kafka.common.utils.Time; import org.apache.kafka.common.utils.Timer; import org.apache.kafka.common.utils.annotation.ApiKeyVersionsSource; @@ -53,7 +56,9 @@ import java.util.List; import java.util.Optional; import java.util.Properties; +import java.util.Random; import java.util.concurrent.BlockingQueue; +import java.util.concurrent.TimeUnit; import static org.apache.kafka.clients.consumer.internals.ConsumerTestBuilder.DEFAULT_GROUP_INSTANCE_ID; import static org.apache.kafka.clients.consumer.internals.ConsumerTestBuilder.DEFAULT_HEARTBEAT_INTERVAL_MS; @@ -62,6 +67,7 @@ import static org.apache.kafka.clients.consumer.internals.ConsumerTestBuilder.DEFAULT_RETRY_BACKOFF_MAX_MS; import static org.apache.kafka.clients.consumer.internals.ConsumerTestBuilder.DEFAULT_RETRY_BACKOFF_MS; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; @@ -79,11 +85,11 @@ public class HeartbeatRequestManagerTest { private int maxPollIntervalMs = DEFAULT_MAX_POLL_INTERVAL_MS; private long retryBackoffMaxMs = DEFAULT_RETRY_BACKOFF_MAX_MS; private static final String DEFAULT_GROUP_ID = "groupId"; + private static final String CONSUMER_COORDINATOR_METRICS = "consumer-coordinator-metrics"; private ConsumerTestBuilder testBuilder; private Time time; private Timer pollTimer; - private ConsumerConfig config; private CoordinatorRequestManager coordinatorRequestManager; private SubscriptionState subscriptions; private Metadata metadata; @@ -95,6 +101,7 @@ public class HeartbeatRequestManagerTest { private final int memberEpoch = 1; private BackgroundEventHandler backgroundEventHandler; private BlockingQueue backgroundEventQueue; + private Metrics metrics; @BeforeEach public void setUp() { @@ -113,7 +120,7 @@ private void setUp(Optional groupInfo) { subscriptions = testBuilder.subscriptions; membershipManager = testBuilder.membershipManager.orElseThrow(IllegalStateException::new); metadata = testBuilder.metadata; - config = testBuilder.config; + metrics = new Metrics(time); when(coordinatorRequestManager.coordinator()).thenReturn(Optional.of(new Node(1, "localhost", 9999))); } @@ -550,8 +557,7 @@ public void testPollTimerExpiration() { when(membershipManager.state()).thenReturn(MemberState.STABLE); time.sleep(maxPollIntervalMs); - NetworkClientDelegate.PollResult pollResult = heartbeatRequestManager.poll(time.milliseconds()); - assertEquals(1, pollResult.unsentRequests.size()); + assertHeartbeat(heartbeatRequestManager, heartbeatIntervalMs); verify(heartbeatState).reset(); verify(heartbeatRequestState).reset(); verify(membershipManager).transitionToStale(); @@ -559,13 +565,60 @@ public void testPollTimerExpiration() { assertNoHeartbeat(heartbeatRequestManager); heartbeatRequestManager.resetPollTimer(time.milliseconds()); assertTrue(pollTimer.notExpired()); - assertHeartbeat(heartbeatRequestManager); + assertHeartbeat(heartbeatRequestManager, heartbeatIntervalMs); + } + + @Test + public void testHeartbeatMetrics() { + // setup + coordinatorRequestManager = mock(CoordinatorRequestManager.class); + membershipManager = mock(MembershipManager.class); + heartbeatState = mock(HeartbeatRequestManager.HeartbeatState.class); + time = new MockTime(); + metrics = new Metrics(time); + heartbeatRequestState = new HeartbeatRequestManager.HeartbeatRequestState( + new LogContext(), + time, + 0, // This initial interval should be 0 to ensure heartbeat on the clock + retryBackoffMs, + retryBackoffMaxMs, + 0); + backgroundEventHandler = mock(BackgroundEventHandler.class); + heartbeatRequestManager = createHeartbeatRequestManager( + coordinatorRequestManager, + membershipManager, + heartbeatState, + heartbeatRequestState, + backgroundEventHandler); + when(coordinatorRequestManager.coordinator()).thenReturn(Optional.of(new Node(1, "localhost", 9999))); + when(membershipManager.state()).thenReturn(MemberState.STABLE); + + assertNotNull(getMetric("heartbeat-response-time-max")); + assertNotNull(getMetric("heartbeat-rate")); + assertNotNull(getMetric("heartbeat-total")); + assertNotNull(getMetric("last-heartbeat-seconds-ago")); + + // test poll + assertHeartbeat(heartbeatRequestManager, 0); + time.sleep(heartbeatIntervalMs); + assertEquals(1.0, getMetric("heartbeat-total").metricValue()); + assertEquals((double) TimeUnit.MILLISECONDS.toSeconds(heartbeatIntervalMs), getMetric("last-heartbeat-seconds-ago").metricValue()); + + assertHeartbeat(heartbeatRequestManager, heartbeatIntervalMs); + assertEquals(0.06d, (double) getMetric("heartbeat-rate").metricValue(), 0.005d); + assertEquals(2.0, getMetric("heartbeat-total").metricValue()); + + // Randomly sleep for some time + Random rand = new Random(); + int randomSleepS = rand.nextInt(11); + time.sleep(randomSleepS * 1000); + assertEquals((double) randomSleepS, getMetric("last-heartbeat-seconds-ago").metricValue()); } - private void assertHeartbeat(HeartbeatRequestManager hrm) { + private void assertHeartbeat(HeartbeatRequestManager hrm, int nextPollMs) { NetworkClientDelegate.PollResult pollResult = hrm.poll(time.milliseconds()); assertEquals(1, pollResult.unsentRequests.size()); - assertEquals(DEFAULT_HEARTBEAT_INTERVAL_MS, pollResult.timeUntilNextPollMs); + assertEquals(nextPollMs, pollResult.timeUntilNextPollMs); pollResult.unsentRequests.get(0).handler().onComplete(createHeartbeatResponse(pollResult.unsentRequests.get(0), Errors.NONE)); } @@ -671,18 +724,8 @@ private ConsumerConfig config() { return new ConsumerConfig(prop); } - private HeartbeatRequestManager createHeartbeatRequestManager() { - LogContext logContext = new LogContext(); - pollTimer = time.timer(maxPollIntervalMs); - return new HeartbeatRequestManager( - logContext, - pollTimer, - config(), - coordinatorRequestManager, - membershipManager, - heartbeatState, - heartbeatRequestState, - backgroundEventHandler); + private KafkaMetric getMetric(final String name) { + return metrics.metrics().get(metrics.metricName(name, CONSUMER_COORDINATOR_METRICS)); } private HeartbeatRequestManager createHeartbeatRequestManager( @@ -701,6 +744,7 @@ private HeartbeatRequestManager createHeartbeatRequestManager( membershipManager, heartbeatState, heartbeatRequestState, - backgroundEventHandler); + backgroundEventHandler, + metrics); } } diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/KafkaConsumerMetricsTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/KafkaConsumerMetricsTest.java index 087f90b7efa31..c75ee906e535b 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/KafkaConsumerMetricsTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/KafkaConsumerMetricsTest.java @@ -17,6 +17,7 @@ package org.apache.kafka.clients.consumer.internals; +import org.apache.kafka.clients.consumer.internals.metrics.KafkaConsumerMetrics; import org.apache.kafka.common.metrics.Metrics; import org.junit.jupiter.api.Test; diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/MembershipManagerImplTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/MembershipManagerImplTest.java index 7683a521e0d6a..a87b6efe16152 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/MembershipManagerImplTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/MembershipManagerImplTest.java @@ -20,7 +20,7 @@ import org.apache.kafka.clients.consumer.internals.events.BackgroundEventHandler; import org.apache.kafka.clients.consumer.internals.events.ConsumerRebalanceListenerCallbackCompletedEvent; import org.apache.kafka.clients.consumer.internals.events.ConsumerRebalanceListenerCallbackNeededEvent; -import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetrics; +import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetricsManager; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.TopicIdPartition; import org.apache.kafka.common.TopicPartition; @@ -1388,7 +1388,7 @@ private ConsumerRebalanceListenerInvoker consumerRebalanceListenerInvoker() { new LogContext(), subscriptionState, new MockTime(1), - new RebalanceCallbackMetrics(new Metrics()) + new RebalanceCallbackMetricsManager(new Metrics()) ); } diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/HeartbeatMetricsManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/HeartbeatMetricsManagerTest.java new file mode 100644 index 0000000000000..1761bc2fef07c --- /dev/null +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/HeartbeatMetricsManagerTest.java @@ -0,0 +1,62 @@ +/* + * 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.clients.consumer.internals.metrics; + +import org.apache.kafka.common.metrics.Metrics; +import org.apache.kafka.common.utils.MockTime; +import org.apache.kafka.common.utils.Time; +import org.junit.jupiter.api.Test; + +import java.util.Random; +import java.util.concurrent.TimeUnit; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; + +public class HeartbeatMetricsManagerTest { + private Time time = new MockTime(); + private Metrics metrics = new Metrics(time); + + @Test + public void testHeartbeatMetrics() { + // Assuming 'metrics' is an instance of your Metrics class + HeartbeatMetricsManager heartbeatMetricsManager = new HeartbeatMetricsManager(metrics); + + // Assert the existence of metrics + assertNotNull(metrics.metric(heartbeatMetricsManager.heartbeatResponseTimeMax)); + assertNotNull(metrics.metric(heartbeatMetricsManager.heartbeatRate)); + assertNotNull(metrics.metric(heartbeatMetricsManager.heartbeatTotal)); + + // Record heartbeat sent time and request latencies + long currentTimeMs = time.milliseconds(); + heartbeatMetricsManager.recordHeartbeatSentMs(currentTimeMs); + heartbeatMetricsManager.recordRequestLatency(100); + heartbeatMetricsManager.recordRequestLatency(103); + heartbeatMetricsManager.recordRequestLatency(102); + + // Assert recorded metrics values + assertEquals(103d, metrics.metric(heartbeatMetricsManager.heartbeatResponseTimeMax).metricValue()); + assertEquals(0.1d, (double) metrics.metric(heartbeatMetricsManager.heartbeatRate).metricValue(), 0.01d); + assertEquals(3d, metrics.metric(heartbeatMetricsManager.heartbeatTotal).metricValue()); + + // Randomly sleep 1-10 seconds + Random rand = new Random(); + int randomSleepS = rand.nextInt(10) + 1; + time.sleep(TimeUnit.SECONDS.toMillis(randomSleepS)); + assertEquals((double) randomSleepS, metrics.metric(heartbeatMetricsManager.lastHeartbeatSecondsAgo).metricValue()); + } +} diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/OffsetCommitMetricsManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/OffsetCommitMetricsManagerTest.java new file mode 100644 index 0000000000000..fa89456019a84 --- /dev/null +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/OffsetCommitMetricsManagerTest.java @@ -0,0 +1,55 @@ +/* + * 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.clients.consumer.internals.metrics; + +import org.apache.kafka.common.metrics.Metrics; +import org.apache.kafka.common.utils.MockTime; +import org.apache.kafka.common.utils.Time; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; + +public class OffsetCommitMetricsManagerTest { + private Time time = new MockTime(); + private Metrics metrics = new Metrics(time); + + @Test + public void testOffsetCommitMetrics() { + // Assuming 'metrics' is an instance of your Metrics class + + // Create an instance of OffsetCommitMetricsManager + OffsetCommitMetricsManager metricsManager = new OffsetCommitMetricsManager(metrics); + + // Assert the existence of metrics + assertNotNull(metrics.metric(metricsManager.commitLatencyAvg)); + assertNotNull(metrics.metric(metricsManager.commitLatencyMax)); + assertNotNull(metrics.metric(metricsManager.commitRate)); + assertNotNull(metrics.metric(metricsManager.commitTotal)); + + // Record request latency + metricsManager.recordRequestLatency(100); + metricsManager.recordRequestLatency(102); + metricsManager.recordRequestLatency(98); + + // Assert the recorded latency + assertEquals(100d, metrics.metric(metricsManager.commitLatencyAvg).metricValue()); + assertEquals(102d, metrics.metric(metricsManager.commitLatencyMax).metricValue()); + assertEquals(0.1d, metrics.metric(metricsManager.commitRate).metricValue()); + assertEquals(3d, metrics.metric(metricsManager.commitTotal).metricValue()); + } +} diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/RebalanceCallbackMetricsManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/RebalanceCallbackMetricsManagerTest.java new file mode 100644 index 0000000000000..c0b5e6eae4824 --- /dev/null +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/RebalanceCallbackMetricsManagerTest.java @@ -0,0 +1,52 @@ +/* + * 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.clients.consumer.internals.metrics; + +import org.apache.kafka.common.metrics.Metrics; +import org.apache.kafka.common.utils.MockTime; +import org.apache.kafka.common.utils.Time; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; + +public class RebalanceCallbackMetricsManagerTest { + private Time time = new MockTime(); + private Metrics metrics = new Metrics(time); + + @Test + public void testRebalanceCallbackMetrics() { + RebalanceCallbackMetricsManager metricsManager = new RebalanceCallbackMetricsManager(metrics); + assertNotNull(metrics.metric(metricsManager.partitionRevokeLatencyAvg)); + assertNotNull(metrics.metric(metricsManager.partitionRevokeLatencyMax)); + assertNotNull(metrics.metric(metricsManager.partitionAssignLatencyAvg)); + assertNotNull(metrics.metric(metricsManager.partitionAssignLatencyMax)); + assertNotNull(metrics.metric(metricsManager.partitionLostLatencyAvg)); + assertNotNull(metrics.metric(metricsManager.partitionLostLatencyMax)); + + metricsManager.recordPartitionsAssignedLatency(100); + metricsManager.recordPartitionsRevokedLatency(101); + metricsManager.recordPartitionsLostLatency(102); + + assertEquals(101d, metrics.metric(metricsManager.partitionRevokeLatencyAvg).metricValue()); + assertEquals(101d, metrics.metric(metricsManager.partitionRevokeLatencyMax).metricValue()); + assertEquals(100d, metrics.metric(metricsManager.partitionAssignLatencyAvg).metricValue()); + assertEquals(100d, metrics.metric(metricsManager.partitionAssignLatencyMax).metricValue()); + assertEquals(102d, metrics.metric(metricsManager.partitionLostLatencyAvg).metricValue()); + assertEquals(102d, metrics.metric(metricsManager.partitionLostLatencyMax).metricValue()); + } +} diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/GroupState.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/GroupState.java deleted file mode 100644 index c160f5acadf4f..0000000000000 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/GroupState.java +++ /dev/null @@ -1,35 +0,0 @@ -/* - * 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.tools.consumer.group; - -import org.apache.kafka.common.Node; - -class GroupState { - public final String group; - public final Node coordinator; - public final String assignmentStrategy; - public final String state; - public final int numMembers; - - public GroupState(String group, Node coordinator, String assignmentStrategy, String state, int numMembers) { - this.group = group; - this.coordinator = coordinator; - this.assignmentStrategy = assignmentStrategy; - this.state = state; - this.numMembers = numMembers; - } -}