From 289bb811573c048317c9b4b4c72c2d623acbb868 Mon Sep 17 00:00:00 2001 From: Nick Telford Date: Fri, 17 May 2024 12:23:49 +0100 Subject: [PATCH 1/2] KAFKA-15541: Add iterator-duration metrics Part of [KIP-989](https://cwiki.apache.org/confluence/x/9KCzDw). This new `StateStore` metric tracks the average and maximum amount of time between creating and closing Iterators. Iterators with very high durations can indicate to users performance problems that should be addressed. If a store reports no data for these metrics, despite the user opening Iterators on the store, it suggests those iterators are not being closed, and have therefore leaked. --- .../state/internals/MeteredKeyValueStore.java | 10 ++++-- ...MeteredMultiVersionedKeyQueryIterator.java | 12 +++++++ .../state/internals/MeteredSessionStore.java | 12 +++++++ .../MeteredTimestampedKeyValueStore.java | 4 ++- .../MeteredVersionedKeyValueStore.java | 8 ++++- .../state/internals/MeteredWindowStore.java | 12 +++++++ .../internals/MeteredWindowStoreIterator.java | 13 ++++--- .../MeteredWindowedKeyValueIterator.java | 13 ++++--- .../internals/metrics/StateStoreMetrics.java | 25 +++++++++++++ .../internals/MeteredKeyValueStoreTest.java | 36 +++++++++++++++++-- .../internals/MeteredSessionStoreTest.java | 34 +++++++++++++++++- .../MeteredTimestampedKeyValueStoreTest.java | 35 ++++++++++++++++-- .../MeteredVersionedKeyValueStoreTest.java | 36 +++++++++++++++++++ .../internals/MeteredWindowStoreTest.java | 34 +++++++++++++++++- 14 files changed, 266 insertions(+), 18 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredKeyValueStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredKeyValueStore.java index 88b49cd21803e..85a7fe9e6bea1 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredKeyValueStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredKeyValueStore.java @@ -90,6 +90,7 @@ public class MeteredKeyValueStore private Sensor prefixScanSensor; private Sensor flushSensor; private Sensor e2eLatencySensor; + protected Sensor iteratorDurationSensor; protected InternalProcessorContext context; private StreamsMetricsImpl streamsMetrics; private TaskId taskId; @@ -165,6 +166,7 @@ private void registerMetrics() { flushSensor = StateStoreMetrics.flushSensor(taskId.toString(), metricsScope, name(), streamsMetrics); deleteSensor = StateStoreMetrics.deleteSensor(taskId.toString(), metricsScope, name(), streamsMetrics); e2eLatencySensor = StateStoreMetrics.e2ELatencySensor(taskId.toString(), metricsScope, name(), streamsMetrics); + iteratorDurationSensor = StateStoreMetrics.iteratorDurationSensor(taskId.toString(), metricsScope, name(), streamsMetrics); StateStoreMetrics.addNumOpenIteratorsGauge(taskId.toString(), metricsScope, name(), streamsMetrics, (config, now) -> numOpenIterators.get()); } @@ -486,7 +488,9 @@ public void close() { try { iter.close(); } finally { - sensor.record(time.nanoseconds() - startNs); + final long duration = time.nanoseconds() - startNs; + sensor.record(duration); + iteratorDurationSensor.record(duration); numOpenIterators.decrementAndGet(); } } @@ -532,7 +536,9 @@ public void close() { try { iter.close(); } finally { - sensor.record(time.nanoseconds() - startNs); + final long duration = time.nanoseconds() - startNs; + sensor.record(duration); + iteratorDurationSensor.record(duration); numOpenIterators.decrementAndGet(); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredMultiVersionedKeyQueryIterator.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredMultiVersionedKeyQueryIterator.java index be695501cafde..4663ef5abbfed 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredMultiVersionedKeyQueryIterator.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredMultiVersionedKeyQueryIterator.java @@ -18,6 +18,9 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Function; + +import org.apache.kafka.common.metrics.Sensor; +import org.apache.kafka.common.utils.Time; import org.apache.kafka.streams.state.VersionedRecordIterator; import org.apache.kafka.streams.state.VersionedRecord; @@ -26,13 +29,21 @@ public class MeteredMultiVersionedKeyQueryIterator implements VersionedRecord private final VersionedRecordIterator iterator; private final Function, VersionedRecord> deserializeValue; private final AtomicInteger numOpenIterators; + private final Sensor sensor; + private final Time time; + private final long startNs; public MeteredMultiVersionedKeyQueryIterator(final VersionedRecordIterator iterator, + final Sensor sensor, + final Time time, final Function, VersionedRecord> deserializeValue, final AtomicInteger numOpenIterators) { this.iterator = iterator; this.deserializeValue = deserializeValue; this.numOpenIterators = numOpenIterators; + this.sensor = sensor; + this.time = time; + this.startNs = time.nanoseconds(); numOpenIterators.incrementAndGet(); } @@ -42,6 +53,7 @@ public void close() { try { iterator.close(); } finally { + sensor.record(time.nanoseconds() - startNs); numOpenIterators.decrementAndGet(); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStore.java index 4bcbb483a314e..0cb9445a92c2f 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStore.java @@ -68,6 +68,7 @@ public class MeteredSessionStore private Sensor flushSensor; private Sensor removeSensor; private Sensor e2eLatencySensor; + private Sensor iteratorDurationSensor; private InternalProcessorContext context; private TaskId taskId; @@ -134,6 +135,7 @@ private void registerMetrics() { flushSensor = StateStoreMetrics.flushSensor(taskId.toString(), metricsScope, name(), streamsMetrics); removeSensor = StateStoreMetrics.removeSensor(taskId.toString(), metricsScope, name(), streamsMetrics); e2eLatencySensor = StateStoreMetrics.e2ELatencySensor(taskId.toString(), metricsScope, name(), streamsMetrics); + iteratorDurationSensor = StateStoreMetrics.iteratorDurationSensor(taskId.toString(), metricsScope, name(), streamsMetrics); StateStoreMetrics.addNumOpenIteratorsGauge(taskId.toString(), metricsScope, name(), streamsMetrics, (config, now) -> numOpenIterators.get()); } @@ -250,6 +252,7 @@ public KeyValueIterator, V> fetch(final K key) { return new MeteredWindowedKeyValueIterator<>( wrapped().fetch(keyBytes(key)), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -263,6 +266,7 @@ public KeyValueIterator, V> backwardFetch(final K key) { return new MeteredWindowedKeyValueIterator<>( wrapped().backwardFetch(keyBytes(key)), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -277,6 +281,7 @@ public KeyValueIterator, V> fetch(final K keyFrom, return new MeteredWindowedKeyValueIterator<>( wrapped().fetch(keyBytes(keyFrom), keyBytes(keyTo)), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -290,6 +295,7 @@ public KeyValueIterator, V> backwardFetch(final K keyFrom, return new MeteredWindowedKeyValueIterator<>( wrapped().backwardFetch(keyBytes(keyFrom), keyBytes(keyTo)), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -310,6 +316,7 @@ public KeyValueIterator, V> findSessions(final K key, earliestSessionEndTime, latestSessionStartTime), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -330,6 +337,7 @@ public KeyValueIterator, V> backwardFindSessions(final K key, latestSessionStartTime ), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -352,6 +360,7 @@ public KeyValueIterator, V> findSessions(final K keyFrom, earliestSessionEndTime, latestSessionStartTime), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -365,6 +374,7 @@ public KeyValueIterator, V> findSessions(final long earliestSessionE return new MeteredWindowedKeyValueIterator<>( wrapped().findSessions(earliestSessionEndTime, latestSessionEndTime), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -387,6 +397,7 @@ public KeyValueIterator, V> backwardFindSessions(final K keyFrom, latestSessionStartTime ), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -458,6 +469,7 @@ private QueryResult runRangeQuery(final Query query, new MeteredWindowedKeyValueIterator<>( rawResult.getResult(), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, StoreQueryUtils.getDeserializeValue(serdes, wrapped()), diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredTimestampedKeyValueStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredTimestampedKeyValueStore.java index d0fcb0cf0ed4b..3c39998389385 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredTimestampedKeyValueStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredTimestampedKeyValueStore.java @@ -350,7 +350,9 @@ public void close() { try { iter.close(); } finally { - sensor.record(time.nanoseconds() - startNs); + final long duration = time.nanoseconds() - startNs; + sensor.record(duration); + iteratorDurationSensor.record(duration); numOpenIterators.decrementAndGet(); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredVersionedKeyValueStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredVersionedKeyValueStore.java index 0c929308d84bc..836bdfb3c4a05 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredVersionedKeyValueStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredVersionedKeyValueStore.java @@ -264,7 +264,13 @@ private QueryResult runMultiVersionedKeyQuery(final Query query, final final QueryResult> rawResult = wrapped().query(rawKeyQuery, positionBound, config); if (rawResult.isSuccess()) { final MeteredMultiVersionedKeyQueryIterator typedResult = - new MeteredMultiVersionedKeyQueryIterator(rawResult.getResult(), StoreQueryUtils.getDeserializeValue(plainValueSerdes), numOpenIterators); + new MeteredMultiVersionedKeyQueryIterator( + rawResult.getResult(), + iteratorDurationSensor, + time, + StoreQueryUtils.getDeserializeValue(plainValueSerdes), + numOpenIterators + ); final QueryResult> typedQueryResult = InternalQueryResultUtil.copyAndSubstituteDeserializedResult(rawResult, typedResult); result = (QueryResult) typedQueryResult; diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredWindowStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredWindowStore.java index 2d63e3ca7c5b3..3f7289e89ac5c 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredWindowStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredWindowStore.java @@ -72,6 +72,7 @@ public class MeteredWindowStore private Sensor fetchSensor; private Sensor flushSensor; private Sensor e2eLatencySensor; + private Sensor iteratorDurationSensor; private InternalProcessorContext context; private TaskId taskId; @@ -153,6 +154,7 @@ private void registerMetrics() { fetchSensor = StateStoreMetrics.fetchSensor(taskId.toString(), metricsScope, name(), streamsMetrics); flushSensor = StateStoreMetrics.flushSensor(taskId.toString(), metricsScope, name(), streamsMetrics); e2eLatencySensor = StateStoreMetrics.e2ELatencySensor(taskId.toString(), metricsScope, name(), streamsMetrics); + iteratorDurationSensor = StateStoreMetrics.iteratorDurationSensor(taskId.toString(), metricsScope, name(), streamsMetrics); StateStoreMetrics.addNumOpenIteratorsGauge(taskId.toString(), metricsScope, name(), streamsMetrics, (config, now) -> numOpenIterators.get()); } @@ -239,6 +241,7 @@ public WindowStoreIterator fetch(final K key, return new MeteredWindowStoreIterator<>( wrapped().fetch(keyBytes(key), timeFrom, timeTo), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::valueFrom, time, @@ -254,6 +257,7 @@ public WindowStoreIterator backwardFetch(final K key, return new MeteredWindowStoreIterator<>( wrapped().backwardFetch(keyBytes(key), timeFrom, timeTo), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::valueFrom, time, @@ -273,6 +277,7 @@ public KeyValueIterator, V> fetch(final K keyFrom, timeFrom, timeTo), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -292,6 +297,7 @@ public KeyValueIterator, V> backwardFetch(final K keyFrom, timeFrom, timeTo), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -305,6 +311,7 @@ public KeyValueIterator, V> fetchAll(final long timeFrom, return new MeteredWindowedKeyValueIterator<>( wrapped().fetchAll(timeFrom, timeTo), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -318,6 +325,7 @@ public KeyValueIterator, V> backwardFetchAll(final long timeFrom, return new MeteredWindowedKeyValueIterator<>( wrapped().backwardFetchAll(timeFrom, timeTo), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -330,6 +338,7 @@ public KeyValueIterator, V> all() { return new MeteredWindowedKeyValueIterator<>( wrapped().all(), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -343,6 +352,7 @@ public KeyValueIterator, V> backwardAll() { return new MeteredWindowedKeyValueIterator<>( wrapped().backwardAll(), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, serdes::valueFrom, @@ -420,6 +430,7 @@ private QueryResult runRangeQuery(final Query query, new MeteredWindowedKeyValueIterator<>( rawResult.getResult(), fetchSensor, + iteratorDurationSensor, streamsMetrics, serdes::keyFrom, getDeserializeValue(serdes, wrapped()), @@ -471,6 +482,7 @@ private QueryResult runKeyQuery(final Query query, final MeteredWindowStoreIterator typedResult = new MeteredWindowStoreIterator<>( rawResult.getResult(), fetchSensor, + iteratorDurationSensor, streamsMetrics, getDeserializeValue(serdes, wrapped()), time, diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredWindowStoreIterator.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredWindowStoreIterator.java index 7f5cf99c0257b..1294cfc1f4504 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredWindowStoreIterator.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredWindowStoreIterator.java @@ -28,7 +28,8 @@ class MeteredWindowStoreIterator implements WindowStoreIterator { private final WindowStoreIterator iter; - private final Sensor sensor; + private final Sensor operationSensor; + private final Sensor iteratorSensor; private final StreamsMetrics metrics; private final Function valueFrom; private final long startNs; @@ -36,13 +37,15 @@ class MeteredWindowStoreIterator implements WindowStoreIterator { private final AtomicInteger numOpenIterators; MeteredWindowStoreIterator(final WindowStoreIterator iter, - final Sensor sensor, + final Sensor operationSensor, + final Sensor iteratorSensor, final StreamsMetrics metrics, final Function valueFrom, final Time time, final AtomicInteger numOpenIterators) { this.iter = iter; - this.sensor = sensor; + this.operationSensor = operationSensor; + this.iteratorSensor = iteratorSensor; this.metrics = metrics; this.valueFrom = valueFrom; this.startNs = time.nanoseconds(); @@ -67,7 +70,9 @@ public void close() { try { iter.close(); } finally { - sensor.record(time.nanoseconds() - startNs); + final long duration = time.nanoseconds() - startNs; + operationSensor.record(duration); + iteratorSensor.record(duration); numOpenIterators.decrementAndGet(); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredWindowedKeyValueIterator.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredWindowedKeyValueIterator.java index a9354c863c81c..e69b27c2c8ed5 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredWindowedKeyValueIterator.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredWindowedKeyValueIterator.java @@ -30,7 +30,8 @@ class MeteredWindowedKeyValueIterator implements KeyValueIterator, V> { private final KeyValueIterator, byte[]> iter; - private final Sensor sensor; + private final Sensor operationSensor; + private final Sensor iteratorSensor; private final StreamsMetrics metrics; private final Function deserializeKey; private final Function deserializeValue; @@ -39,14 +40,16 @@ class MeteredWindowedKeyValueIterator implements KeyValueIterator, byte[]> iter, - final Sensor sensor, + final Sensor operationSensor, + final Sensor iteratorSensor, final StreamsMetrics metrics, final Function deserializeKey, final Function deserializeValue, final Time time, final AtomicInteger numOpenIterators) { this.iter = iter; - this.sensor = sensor; + this.operationSensor = operationSensor; + this.iteratorSensor = iteratorSensor; this.metrics = metrics; this.deserializeKey = deserializeKey; this.deserializeValue = deserializeValue; @@ -77,7 +80,9 @@ public void close() { try { iter.close(); } finally { - sensor.record(time.nanoseconds() - startNs); + final long duration = time.nanoseconds() - startNs; + operationSensor.record(duration); + iteratorSensor.record(duration); numOpenIterators.decrementAndGet(); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/metrics/StateStoreMetrics.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/metrics/StateStoreMetrics.java index eea7a9644221c..58fde290bc1cb 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/metrics/StateStoreMetrics.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/metrics/StateStoreMetrics.java @@ -149,6 +149,14 @@ private StateStoreMetrics() {} private static final String NUM_OPEN_ITERATORS_DESCRIPTION = "The current number of iterators on the store that have been created, but not yet closed"; + private static final String ITERATOR_DURATION = "iterator-duration"; + private static final String ITERATOR_DURATION_DESCRIPTION = + "time spent between creating an Iterator and closing it, in nanoseconds"; + private static final String ITERATOR_DURATION_AVG_DESCRIPTION = + AVG_DESCRIPTION_PREFIX + ITERATOR_DURATION_DESCRIPTION; + private static final String ITERATOR_DURATION_MAX_DESCRIPTION = + MAX_DESCRIPTION_PREFIX + ITERATOR_DURATION_DESCRIPTION; + public static Sensor putSensor(final String taskId, final String storeType, final String storeName, @@ -409,6 +417,23 @@ public static Sensor e2ELatencySensor(final String taskId, return sensor; } + public static Sensor iteratorDurationSensor(final String taskId, + final String storeType, + final String storeName, + final StreamsMetricsImpl streamsMetrics) { + final Sensor sensor = streamsMetrics.storeLevelSensor(taskId, storeName, ITERATOR_DURATION, RecordingLevel.DEBUG); + final Map tagMap = streamsMetrics.storeLevelTagMap(taskId, storeType, storeName); + addAvgAndMaxToSensor( + sensor, + STATE_STORE_LEVEL_GROUP, + tagMap, + ITERATOR_DURATION, + ITERATOR_DURATION_AVG_DESCRIPTION, + ITERATOR_DURATION_MAX_DESCRIPTION + ); + return sensor; + } + public static void addNumOpenIteratorsGauge(final String taskId, final String storeType, final String storeName, diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredKeyValueStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredKeyValueStoreTest.java index 00151c6798105..d2227bc69b103 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredKeyValueStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredKeyValueStoreTest.java @@ -30,7 +30,6 @@ import org.apache.kafka.common.serialization.StringSerializer; import org.apache.kafka.common.utils.Bytes; import org.apache.kafka.common.utils.MockTime; -import org.apache.kafka.common.utils.Time; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.processor.ProcessorContext; @@ -51,6 +50,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import static org.apache.kafka.common.utils.Utils.mkEntry; @@ -98,10 +98,12 @@ public class MeteredKeyValueStoreTest { private MeteredKeyValueStore metered; private final Metrics metrics = new Metrics(); private Map tags; + private MockTime mockTime; @Before public void before() { - final Time mockTime = new MockTime(); + final MockTime mockTime = new MockTime(); + this.mockTime = mockTime; metered = new MeteredKeyValueStore<>( inner, STORE_TYPE, @@ -458,6 +460,36 @@ public void shouldTrackOpenIteratorsMetric() { assertThat((Integer) openIteratorsMetric.metricValue(), equalTo(0)); } + @Test + public void shouldTimeIteratorDuration() { + when(inner.all()).thenReturn(KeyValueIterators.emptyIterator()); + init(); + + final KafkaMetric iteratorDurationAvgMetric = metric("iterator-duration-avg"); + final KafkaMetric iteratorDurationMaxMetric = metric("iterator-duration-max"); + assertThat(iteratorDurationAvgMetric, not(nullValue())); + assertThat(iteratorDurationMaxMetric, not(nullValue())); + + assertThat((Double) iteratorDurationAvgMetric.metricValue(), equalTo(Double.NaN)); + assertThat((Double) iteratorDurationMaxMetric.metricValue(), equalTo(Double.NaN)); + + try (final KeyValueIterator iterator = metered.all()) { + // nothing to do, just close immediately + mockTime.sleep(2); + } + + assertThat((double) iteratorDurationAvgMetric.metricValue(), equalTo(2.0 * TimeUnit.MILLISECONDS.toNanos(1))); + assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(2.0 * TimeUnit.MILLISECONDS.toNanos(1))); + + try (final KeyValueIterator iterator = metered.all()) { + // nothing to do, just close immediately + mockTime.sleep(3); + } + + assertThat((double) iteratorDurationAvgMetric.metricValue(), equalTo(2.5 * TimeUnit.MILLISECONDS.toNanos(1))); + assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(3.0 * TimeUnit.MILLISECONDS.toNanos(1))); + } + private KafkaMetric metric(final MetricName metricName) { return this.metrics.metric(metricName); } diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreTest.java index b7c99032f8920..3c5a923e6d42b 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreTest.java @@ -53,6 +53,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import static org.apache.kafka.common.utils.Utils.mkEntry; @@ -95,6 +96,7 @@ public class MeteredSessionStoreTest { private final String threadId = Thread.currentThread().getName(); private final TaskId taskId = new TaskId(0, 0, "My-Topology"); private final Metrics metrics = new Metrics(); + private MockTime mockTime; private MeteredSessionStore store; @Mock private SessionStore innerStore; @@ -105,7 +107,7 @@ public class MeteredSessionStoreTest { @Before public void before() { - final Time mockTime = new MockTime(); + mockTime = new MockTime(); store = new MeteredSessionStore<>( innerStore, STORE_TYPE, @@ -622,6 +624,36 @@ public void shouldTrackOpenIteratorsMetric() { assertThat((Integer) openIteratorsMetric.metricValue(), equalTo(0)); } + @Test + public void shouldTimeIteratorDuration() { + when(innerStore.backwardFetch(KEY_BYTES)).thenReturn(KeyValueIterators.emptyIterator()); + init(); + + final KafkaMetric iteratorDurationAvgMetric = metric("iterator-duration-avg"); + final KafkaMetric iteratorDurationMaxMetric = metric("iterator-duration-max"); + assertThat(iteratorDurationAvgMetric, not(nullValue())); + assertThat(iteratorDurationMaxMetric, not(nullValue())); + + assertThat((Double) iteratorDurationAvgMetric.metricValue(), equalTo(Double.NaN)); + assertThat((Double) iteratorDurationMaxMetric.metricValue(), equalTo(Double.NaN)); + + try (final KeyValueIterator, String> iterator = store.backwardFetch(KEY)) { + // nothing to do, just close immediately + mockTime.sleep(2); + } + + assertThat((double) iteratorDurationAvgMetric.metricValue(), equalTo(2.0 * TimeUnit.MILLISECONDS.toNanos(1))); + assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(2.0 * TimeUnit.MILLISECONDS.toNanos(1))); + + try (final KeyValueIterator, String> iterator = store.backwardFetch(KEY)) { + // nothing to do, just close immediately + mockTime.sleep(3); + } + + assertThat((double) iteratorDurationAvgMetric.metricValue(), equalTo(2.5 * TimeUnit.MILLISECONDS.toNanos(1))); + assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(3.0 * TimeUnit.MILLISECONDS.toNanos(1))); + } + private KafkaMetric metric(final String name) { return this.metrics.metric(new MetricName(name, STORE_LEVEL_GROUP, "", this.tags)); } diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedKeyValueStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedKeyValueStoreTest.java index f629acfc9a8a6..7e37440a075bd 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedKeyValueStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedKeyValueStoreTest.java @@ -29,7 +29,6 @@ import org.apache.kafka.common.serialization.Serializer; import org.apache.kafka.common.utils.Bytes; import org.apache.kafka.common.utils.MockTime; -import org.apache.kafka.common.utils.Time; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.errors.StreamsException; @@ -53,6 +52,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import static org.apache.kafka.common.utils.Utils.mkEntry; import static org.apache.kafka.common.utils.Utils.mkMap; @@ -95,6 +95,7 @@ public class MeteredTimestampedKeyValueStoreTest { private KeyValueStore inner; @Mock private InternalProcessorContext context; + private MockTime mockTime; private final static Map CONFIGS = mkMap(mkEntry(StreamsConfig.InternalConfig.TOPIC_PREFIX_ALTERNATIVE, APPLICATION_ID)); @@ -107,7 +108,7 @@ public class MeteredTimestampedKeyValueStoreTest { @Before public void before() { - final Time mockTime = new MockTime(); + mockTime = new MockTime(); metered = new MeteredTimestampedKeyValueStore<>( inner, "scope", @@ -456,4 +457,34 @@ public void shouldTrackOpenIteratorsMetric() { assertThat((Integer) openIteratorsMetric.metricValue(), equalTo(0)); } + + @Test + public void shouldTimeIteratorDuration() { + when(inner.all()).thenReturn(KeyValueIterators.emptyIterator()); + init(); + + final KafkaMetric iteratorDurationAvgMetric = metric("iterator-duration-avg"); + final KafkaMetric iteratorDurationMaxMetric = metric("iterator-duration-max"); + assertThat(iteratorDurationAvgMetric, not(nullValue())); + assertThat(iteratorDurationMaxMetric, not(nullValue())); + + assertThat((Double) iteratorDurationAvgMetric.metricValue(), equalTo(Double.NaN)); + assertThat((Double) iteratorDurationMaxMetric.metricValue(), equalTo(Double.NaN)); + + try (final KeyValueIterator> iterator = metered.all()) { + // nothing to do, just close immediately + mockTime.sleep(2); + } + + assertThat((double) iteratorDurationAvgMetric.metricValue(), equalTo(2.0 * TimeUnit.MILLISECONDS.toNanos(1))); + assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(2.0 * TimeUnit.MILLISECONDS.toNanos(1))); + + try (final KeyValueIterator> iterator = metered.all()) { + // nothing to do, just close immediately + mockTime.sleep(3); + } + + assertThat((double) iteratorDurationAvgMetric.metricValue(), equalTo(2.5 * TimeUnit.MILLISECONDS.toNanos(1))); + assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(3.0 * TimeUnit.MILLISECONDS.toNanos(1))); + } } diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredVersionedKeyValueStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredVersionedKeyValueStoreTest.java index 99f76839041c7..10b8e2ec6fb1c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredVersionedKeyValueStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredVersionedKeyValueStoreTest.java @@ -39,6 +39,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import org.apache.kafka.common.MetricName; import org.apache.kafka.common.metrics.KafkaMetric; @@ -390,6 +391,41 @@ public void shouldTrackOpenIteratorsMetric() { assertThat((Integer) openIteratorsMetric.metricValue(), equalTo(0)); } + @Test + public void shouldTimeIteratorDuration() { + final MultiVersionedKeyQuery query = MultiVersionedKeyQuery.withKey(KEY); + final PositionBound bound = PositionBound.unbounded(); + final QueryConfig config = new QueryConfig(false); + when(inner.query(any(), any(), any())).thenReturn( + QueryResult.forResult(new LogicalSegmentIterator(Collections.emptyListIterator(), RAW_KEY, 0L, 0L, ResultOrder.ANY))); + + final KafkaMetric iteratorDurationAvgMetric = getMetric("iterator-duration-avg"); + final KafkaMetric iteratorDurationMaxMetric = getMetric("iterator-duration-max"); + assertThat(iteratorDurationAvgMetric, not(nullValue())); + assertThat(iteratorDurationMaxMetric, not(nullValue())); + + assertThat((Double) iteratorDurationAvgMetric.metricValue(), equalTo(Double.NaN)); + assertThat((Double) iteratorDurationMaxMetric.metricValue(), equalTo(Double.NaN)); + + final QueryResult> first = store.query(query, bound, config); + try (final VersionedRecordIterator iterator = first.getResult()) { + // nothing to do, just close immediately + mockTime.sleep(2); + } + + assertThat((double) iteratorDurationAvgMetric.metricValue(), equalTo(2.0 * TimeUnit.MILLISECONDS.toNanos(1))); + assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(2.0 * TimeUnit.MILLISECONDS.toNanos(1))); + + final QueryResult> second = store.query(query, bound, config); + try (final VersionedRecordIterator iterator = second.getResult()) { + // nothing to do, just close immediately + mockTime.sleep(3); + } + + assertThat((double) iteratorDurationAvgMetric.metricValue(), equalTo(2.5 * TimeUnit.MILLISECONDS.toNanos(1))); + assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(3.0 * TimeUnit.MILLISECONDS.toNanos(1))); + } + private KafkaMetric getMetric(final String name) { return metrics.metric(new MetricName(name, STORE_LEVEL_GROUP, "", tags)); } diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredWindowStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredWindowStoreTest.java index 3997a5e549ec1..aa637035c7756 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredWindowStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredWindowStoreTest.java @@ -52,6 +52,7 @@ import java.time.temporal.ChronoUnit; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import static java.time.Instant.ofEpochMilli; @@ -94,11 +95,12 @@ public class MeteredWindowStoreTest { private InternalMockProcessorContext context; @SuppressWarnings("unchecked") private final WindowStore innerStoreMock = mock(WindowStore.class); + private final MockTime mockTime = new MockTime(); private MeteredWindowStore store = new MeteredWindowStore<>( innerStoreMock, WINDOW_SIZE_MS, // any size STORE_TYPE, - new MockTime(), + mockTime, Serdes.String(), new SerdeThatDoesntHandleNull() ); @@ -463,6 +465,36 @@ public void shouldTrackOpenIteratorsMetric() { assertThat((Integer) openIteratorsMetric.metricValue(), equalTo(0)); } + @Test + public void shouldTimeIteratorDuration() { + when(innerStoreMock.all()).thenReturn(KeyValueIterators.emptyIterator()); + store.init((StateStoreContext) context, store); + + final KafkaMetric iteratorDurationAvgMetric = metric("iterator-duration-avg"); + final KafkaMetric iteratorDurationMaxMetric = metric("iterator-duration-max"); + assertThat(iteratorDurationAvgMetric, not(nullValue())); + assertThat(iteratorDurationMaxMetric, not(nullValue())); + + assertThat((Double) iteratorDurationAvgMetric.metricValue(), equalTo(Double.NaN)); + assertThat((Double) iteratorDurationMaxMetric.metricValue(), equalTo(Double.NaN)); + + try (final KeyValueIterator, String> iterator = store.all()) { + // nothing to do, just close immediately + mockTime.sleep(2); + } + + assertThat((double) iteratorDurationAvgMetric.metricValue(), equalTo(2.0 * TimeUnit.MILLISECONDS.toNanos(1))); + assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(2.0 * TimeUnit.MILLISECONDS.toNanos(1))); + + try (final KeyValueIterator, String> iterator = store.all()) { + // nothing to do, just close immediately + mockTime.sleep(3); + } + + assertThat((double) iteratorDurationAvgMetric.metricValue(), equalTo(2.5 * TimeUnit.MILLISECONDS.toNanos(1))); + assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(3.0 * TimeUnit.MILLISECONDS.toNanos(1))); + } + private KafkaMetric metric(final String name) { return metrics.metric(new MetricName(name, STORE_LEVEL_GROUP, "", tags)); } From d69544c4e59cee7b5c5cb0a8e4c8f43a7b8636f3 Mon Sep 17 00:00:00 2001 From: "Matthias J. Sax" Date: Wed, 22 May 2024 21:33:52 -0700 Subject: [PATCH 2/2] Update streams/src/main/java/org/apache/kafka/streams/state/internals/metrics/StateStoreMetrics.java --- .../streams/state/internals/metrics/StateStoreMetrics.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/metrics/StateStoreMetrics.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/metrics/StateStoreMetrics.java index 58fde290bc1cb..aa9d9c3238da6 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/metrics/StateStoreMetrics.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/metrics/StateStoreMetrics.java @@ -151,7 +151,7 @@ private StateStoreMetrics() {} private static final String ITERATOR_DURATION = "iterator-duration"; private static final String ITERATOR_DURATION_DESCRIPTION = - "time spent between creating an Iterator and closing it, in nanoseconds"; + "time spent between creating an iterator and closing it, in nanoseconds"; private static final String ITERATOR_DURATION_AVG_DESCRIPTION = AVG_DESCRIPTION_PREFIX + ITERATOR_DURATION_DESCRIPTION; private static final String ITERATOR_DURATION_MAX_DESCRIPTION =