diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredIterator.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredIterator.java new file mode 100644 index 0000000000000..c04fb43020c53 --- /dev/null +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredIterator.java @@ -0,0 +1,30 @@ +/* + * 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.streams.state.internals; + +/** + * Common super-interface of all Metered Iterator types. + * + * This enables tracking the timestamp the Iterator was first created, for the oldest-iterator-open-since-ms metric. + */ +public interface MeteredIterator { + + /** + * @return The UNIX timestamp, in milliseconds, that this Iterator was created/opened. + */ + long startTimestamp(); +} 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 f4c2b3b6a9903..fbe42b87065e8 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 @@ -49,9 +49,12 @@ import org.apache.kafka.streams.state.internals.metrics.StateStoreMetrics; import java.util.ArrayList; +import java.util.Comparator; import java.util.List; import java.util.Map; +import java.util.NavigableSet; import java.util.Objects; +import java.util.concurrent.ConcurrentSkipListSet; import java.util.concurrent.atomic.LongAdder; import java.util.function.Function; @@ -96,6 +99,7 @@ public class MeteredKeyValueStore private TaskId taskId; protected LongAdder numOpenIterators = new LongAdder(); + protected NavigableSet openIterators = new ConcurrentSkipListSet<>(Comparator.comparingLong(MeteredIterator::startTimestamp)); @SuppressWarnings("rawtypes") private final Map queryHandlers = @@ -169,6 +173,9 @@ private void registerMetrics() { iteratorDurationSensor = StateStoreMetrics.iteratorDurationSensor(taskId.toString(), metricsScope, name(), streamsMetrics); StateStoreMetrics.addNumOpenIteratorsGauge(taskId.toString(), metricsScope, name(), streamsMetrics, (config, now) -> numOpenIterators.sum()); + StateStoreMetrics.addOldestOpenIteratorGauge(taskId.toString(), metricsScope, name(), streamsMetrics, + (config, now) -> openIterators.isEmpty() ? null : openIterators.first().startTimestamp() + ); } protected Serde prepareValueSerdeForStore(final Serde valueSerde, final SerdeGetter getter) { @@ -456,18 +463,26 @@ protected void maybeRecordE2ELatency() { } } - private class MeteredKeyValueIterator implements KeyValueIterator { + private class MeteredKeyValueIterator implements KeyValueIterator, MeteredIterator { private final KeyValueIterator iter; private final Sensor sensor; private final long startNs; + private final long startTimestamp; private MeteredKeyValueIterator(final KeyValueIterator iter, final Sensor sensor) { this.iter = iter; this.sensor = sensor; + this.startTimestamp = time.milliseconds(); this.startNs = time.nanoseconds(); numOpenIterators.increment(); + openIterators.add(this); + } + + @Override + public long startTimestamp() { + return startTimestamp; } @Override @@ -492,6 +507,7 @@ public void close() { sensor.record(duration); iteratorDurationSensor.record(duration); numOpenIterators.decrement(); + openIterators.remove(this); } } @@ -501,11 +517,12 @@ public K peekNextKey() { } } - private class MeteredKeyValueTimestampedIterator implements KeyValueIterator { + private class MeteredKeyValueTimestampedIterator implements KeyValueIterator, MeteredIterator { private final KeyValueIterator iter; private final Sensor sensor; private final long startNs; + private final long startTimestamp; private final Function valueDeserializer; private MeteredKeyValueTimestampedIterator(final KeyValueIterator iter, @@ -514,8 +531,15 @@ private MeteredKeyValueTimestampedIterator(final KeyValueIterator this.iter = iter; this.sensor = sensor; this.valueDeserializer = valueDeserializer; + this.startTimestamp = time.milliseconds(); this.startNs = time.nanoseconds(); numOpenIterators.increment(); + openIterators.add(this); + } + + @Override + public long startTimestamp() { + return startTimestamp; } @Override @@ -540,6 +564,7 @@ public void close() { sensor.record(duration); iteratorDurationSensor.record(duration); numOpenIterators.decrement(); + openIterators.remove(this); } } 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 92c315c09be99..48347365756c9 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 @@ -17,6 +17,7 @@ package org.apache.kafka.streams.state.internals; import java.util.concurrent.atomic.LongAdder; +import java.util.Set; import java.util.function.Function; import org.apache.kafka.common.metrics.Sensor; @@ -24,7 +25,7 @@ import org.apache.kafka.streams.state.VersionedRecordIterator; import org.apache.kafka.streams.state.VersionedRecord; -public class MeteredMultiVersionedKeyQueryIterator implements VersionedRecordIterator { +class MeteredMultiVersionedKeyQueryIterator implements VersionedRecordIterator, MeteredIterator { private final VersionedRecordIterator iterator; private final Function, VersionedRecord> deserializeValue; @@ -32,21 +33,31 @@ public class MeteredMultiVersionedKeyQueryIterator implements VersionedRecord private final Sensor sensor; private final Time time; private final long startNs; + private final long startTimestampMs; + private final Set openIterators; public MeteredMultiVersionedKeyQueryIterator(final VersionedRecordIterator iterator, final Sensor sensor, final Time time, final Function, VersionedRecord> deserializeValue, - final LongAdder numOpenIterators) { + final LongAdder numOpenIterators, + final Set openIterators) { this.iterator = iterator; this.deserializeValue = deserializeValue; this.numOpenIterators = numOpenIterators; + this.openIterators = openIterators; this.sensor = sensor; this.time = time; this.startNs = time.nanoseconds(); + this.startTimestampMs = time.milliseconds(); numOpenIterators.increment(); + openIterators.add(this); } + @Override + public long startTimestamp() { + return startTimestampMs; + } @Override public void close() { @@ -55,6 +66,7 @@ public void close() { } finally { sensor.record(time.nanoseconds() - startNs); numOpenIterators.decrement(); + openIterators.remove(this); } } 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 5208a0cf671c9..731bc3145c181 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 @@ -45,8 +45,11 @@ import org.apache.kafka.streams.state.internals.StoreQueryUtils.QueryHandler; import org.apache.kafka.streams.state.internals.metrics.StateStoreMetrics; +import java.util.Comparator; import java.util.Map; +import java.util.NavigableSet; import java.util.Objects; +import java.util.concurrent.ConcurrentSkipListSet; import java.util.concurrent.atomic.LongAdder; import static org.apache.kafka.common.utils.Utils.mkEntry; @@ -73,6 +76,7 @@ public class MeteredSessionStore private TaskId taskId; private LongAdder numOpenIterators = new LongAdder(); + private final NavigableSet openIterators = new ConcurrentSkipListSet<>(Comparator.comparingLong(MeteredIterator::startTimestamp)); @SuppressWarnings("rawtypes") private final Map queryHandlers = @@ -138,6 +142,9 @@ private void registerMetrics() { iteratorDurationSensor = StateStoreMetrics.iteratorDurationSensor(taskId.toString(), metricsScope, name(), streamsMetrics); StateStoreMetrics.addNumOpenIteratorsGauge(taskId.toString(), metricsScope, name(), streamsMetrics, (config, now) -> numOpenIterators.sum()); + StateStoreMetrics.addOldestOpenIteratorGauge(taskId.toString(), metricsScope, name(), streamsMetrics, + (config, now) -> openIterators.isEmpty() ? null : openIterators.first().startTimestamp() + ); } @@ -257,7 +264,8 @@ public KeyValueIterator, V> fetch(final K key) { serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators); + numOpenIterators, + openIterators); } @Override @@ -271,7 +279,8 @@ public KeyValueIterator, V> backwardFetch(final K key) { serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators + numOpenIterators, + openIterators ); } @@ -286,7 +295,8 @@ public KeyValueIterator, V> fetch(final K keyFrom, serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators); + numOpenIterators, + openIterators); } @Override @@ -300,7 +310,8 @@ public KeyValueIterator, V> backwardFetch(final K keyFrom, serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators + numOpenIterators, + openIterators ); } @@ -321,7 +332,8 @@ public KeyValueIterator, V> findSessions(final K key, serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators); + numOpenIterators, + openIterators); } @Override @@ -342,7 +354,8 @@ public KeyValueIterator, V> backwardFindSessions(final K key, serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators + numOpenIterators, + openIterators ); } @@ -365,7 +378,8 @@ public KeyValueIterator, V> findSessions(final K keyFrom, serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators); + numOpenIterators, + openIterators); } @Override @@ -379,7 +393,8 @@ public KeyValueIterator, V> findSessions(final long earliestSessionE serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators); + numOpenIterators, + openIterators); } @Override @@ -402,7 +417,8 @@ public KeyValueIterator, V> backwardFindSessions(final K keyFrom, serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators + numOpenIterators, + openIterators ); } @@ -474,7 +490,8 @@ private QueryResult runRangeQuery(final Query query, serdes::keyFrom, StoreQueryUtils.getDeserializeValue(serdes, wrapped()), time, - numOpenIterators + numOpenIterators, + openIterators ); final QueryResult> typedQueryResult = InternalQueryResultUtil.copyAndSubstituteDeserializedResult(rawResult, typedResult); 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 757a1bd868cf0..0b4702b9dbea5 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 @@ -307,11 +307,12 @@ private QueryResult runRangeQuery(final Query query, } @SuppressWarnings("unchecked") - private class MeteredTimestampedKeyValueStoreIterator implements KeyValueIterator { + private class MeteredTimestampedKeyValueStoreIterator implements KeyValueIterator, MeteredIterator { private final KeyValueIterator iter; private final Sensor sensor; private final long startNs; + private final long startTimestampMs; private final Function> valueAndTimestampDeserializer; private final boolean returnPlainValue; @@ -324,8 +325,15 @@ private MeteredTimestampedKeyValueStoreIterator(final KeyValueIterator QueryResult runMultiVersionedKeyQuery(final Query query, final iteratorDurationSensor, time, StoreQueryUtils.getDeserializeValue(plainValueSerdes), - numOpenIterators + numOpenIterators, + openIterators ); final QueryResult> typedQueryResult = InternalQueryResultUtil.copyAndSubstituteDeserializedResult(rawResult, typedResult); 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 ce7f0229a3cf1..a62e8c47563f2 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 @@ -47,8 +47,11 @@ import org.apache.kafka.streams.state.internals.StoreQueryUtils.QueryHandler; import org.apache.kafka.streams.state.internals.metrics.StateStoreMetrics; +import java.util.Comparator; import java.util.Map; +import java.util.NavigableSet; import java.util.Objects; +import java.util.concurrent.ConcurrentSkipListSet; import java.util.concurrent.atomic.LongAdder; import static org.apache.kafka.common.utils.Utils.mkEntry; @@ -77,6 +80,7 @@ public class MeteredWindowStore private TaskId taskId; private LongAdder numOpenIterators = new LongAdder(); + private NavigableSet openIterators = new ConcurrentSkipListSet<>(Comparator.comparingLong(MeteredIterator::startTimestamp)); @SuppressWarnings("rawtypes") private final Map queryHandlers = @@ -157,6 +161,9 @@ private void registerMetrics() { iteratorDurationSensor = StateStoreMetrics.iteratorDurationSensor(taskId.toString(), metricsScope, name(), streamsMetrics); StateStoreMetrics.addNumOpenIteratorsGauge(taskId.toString(), metricsScope, name(), streamsMetrics, (config, now) -> numOpenIterators.sum()); + StateStoreMetrics.addOldestOpenIteratorGauge(taskId.toString(), metricsScope, name(), streamsMetrics, + (config, now) -> openIterators.isEmpty() ? null : openIterators.first().startTimestamp() + ); } @Deprecated @@ -245,7 +252,8 @@ public WindowStoreIterator fetch(final K key, streamsMetrics, serdes::valueFrom, time, - numOpenIterators + numOpenIterators, + openIterators ); } @@ -261,7 +269,8 @@ public WindowStoreIterator backwardFetch(final K key, streamsMetrics, serdes::valueFrom, time, - numOpenIterators + numOpenIterators, + openIterators ); } @@ -282,7 +291,8 @@ public KeyValueIterator, V> fetch(final K keyFrom, serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators); + numOpenIterators, + openIterators); } @Override @@ -302,7 +312,8 @@ public KeyValueIterator, V> backwardFetch(final K keyFrom, serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators); + numOpenIterators, + openIterators); } @Override @@ -316,7 +327,8 @@ public KeyValueIterator, V> fetchAll(final long timeFrom, serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators); + numOpenIterators, + openIterators); } @Override @@ -330,7 +342,8 @@ public KeyValueIterator, V> backwardFetchAll(final long timeFrom, serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators); + numOpenIterators, + openIterators); } @Override @@ -343,7 +356,8 @@ public KeyValueIterator, V> all() { serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators + numOpenIterators, + openIterators ); } @@ -357,7 +371,8 @@ public KeyValueIterator, V> backwardAll() { serdes::keyFrom, serdes::valueFrom, time, - numOpenIterators + numOpenIterators, + openIterators ); } @@ -435,7 +450,8 @@ private QueryResult runRangeQuery(final Query query, serdes::keyFrom, getDeserializeValue(serdes, wrapped()), time, - numOpenIterators + numOpenIterators, + openIterators ); final QueryResult> typedQueryResult = InternalQueryResultUtil.copyAndSubstituteDeserializedResult(rawResult, typedResult); @@ -486,7 +502,8 @@ private QueryResult runKeyQuery(final Query query, streamsMetrics, getDeserializeValue(serdes, wrapped()), time, - numOpenIterators + numOpenIterators, + openIterators ); final QueryResult> typedQueryResult = InternalQueryResultUtil.copyAndSubstituteDeserializedResult(rawResult, typedResult); 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 62fb7b1b24867..90bd3e5f96f67 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 @@ -23,9 +23,10 @@ import org.apache.kafka.streams.state.WindowStoreIterator; import java.util.concurrent.atomic.LongAdder; +import java.util.Set; import java.util.function.Function; -class MeteredWindowStoreIterator implements WindowStoreIterator { +class MeteredWindowStoreIterator implements WindowStoreIterator, MeteredIterator { private final WindowStoreIterator iter; private final Sensor operationSensor; @@ -33,8 +34,10 @@ class MeteredWindowStoreIterator implements WindowStoreIterator { private final StreamsMetrics metrics; private final Function valueFrom; private final long startNs; + private final long startTimestampMs; private final Time time; private final LongAdder numOpenIterators; + private final Set openIterators; MeteredWindowStoreIterator(final WindowStoreIterator iter, final Sensor operationSensor, @@ -42,16 +45,25 @@ class MeteredWindowStoreIterator implements WindowStoreIterator { final StreamsMetrics metrics, final Function valueFrom, final Time time, - final LongAdder numOpenIterators) { + final LongAdder numOpenIterators, + final Set openIterators) { this.iter = iter; this.operationSensor = operationSensor; this.iteratorSensor = iteratorSensor; this.metrics = metrics; this.valueFrom = valueFrom; this.startNs = time.nanoseconds(); + this.startTimestampMs = time.milliseconds(); this.time = time; this.numOpenIterators = numOpenIterators; + this.openIterators = openIterators; numOpenIterators.increment(); + openIterators.add(this); + } + + @Override + public long startTimestamp() { + return startTimestampMs; } @Override @@ -74,6 +86,7 @@ public void close() { operationSensor.record(duration); iteratorSensor.record(duration); numOpenIterators.decrement(); + openIterators.remove(this); } } 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 e663a2b3a2036..1eeacd81babc0 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 @@ -25,9 +25,10 @@ import org.apache.kafka.streams.state.KeyValueIterator; import java.util.concurrent.atomic.LongAdder; +import java.util.Set; import java.util.function.Function; -class MeteredWindowedKeyValueIterator implements KeyValueIterator, V> { +class MeteredWindowedKeyValueIterator implements KeyValueIterator, V>, MeteredIterator { private final KeyValueIterator, byte[]> iter; private final Sensor operationSensor; @@ -36,8 +37,10 @@ class MeteredWindowedKeyValueIterator implements KeyValueIterator deserializeKey; private final Function deserializeValue; private final long startNs; + private final long startTimestampMs; private final Time time; private final LongAdder numOpenIterators; + private final Set openIterators; MeteredWindowedKeyValueIterator(final KeyValueIterator, byte[]> iter, final Sensor operationSensor, @@ -46,7 +49,8 @@ class MeteredWindowedKeyValueIterator implements KeyValueIterator deserializeKey, final Function deserializeValue, final Time time, - final LongAdder numOpenIterators) { + final LongAdder numOpenIterators, + final Set openIterators) { this.iter = iter; this.operationSensor = operationSensor; this.iteratorSensor = iteratorSensor; @@ -54,9 +58,17 @@ class MeteredWindowedKeyValueIterator implements KeyValueIterator oldestOpenIteratorGauge) { + streamsMetrics.addStoreLevelMutableMetric( + taskId, + storeType, + storeName, + OLDEST_ITERATOR_OPEN_SINCE_MS, + OLDEST_ITERATOR_OPEN_SINCE_MS_DESCRIPTION, + RecordingLevel.INFO, + oldestOpenIteratorGauge + ); + } + private static Sensor sizeOrCountSensor(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 7e78084027e2c..860be7f7efe04 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 @@ -490,6 +490,42 @@ public void shouldTimeIteratorDuration() { assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(3.0 * TimeUnit.MILLISECONDS.toNanos(1))); } + @Test + public void shouldTrackOldestOpenIteratorTimestamp() { + when(inner.all()).thenReturn(KeyValueIterators.emptyIterator()); + init(); + + final KafkaMetric oldestIteratorTimestampMetric = metric("oldest-iterator-open-since-ms"); + assertThat(oldestIteratorTimestampMetric, not(nullValue())); + + assertThat(oldestIteratorTimestampMetric.metricValue(), nullValue()); + + KeyValueIterator second = null; + final long secondTimestamp; + try { + try (final KeyValueIterator first = metered.all()) { + final long oldestTimestamp = mockTime.milliseconds(); + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(oldestTimestamp)); + mockTime.sleep(100); + + // open a second iterator before closing the first to test that we still produce the first iterator's timestamp + second = metered.all(); + secondTimestamp = mockTime.milliseconds(); + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(oldestTimestamp)); + mockTime.sleep(100); + } + + // now that the first iterator is closed, check that the timestamp has advanced to the still open second iterator + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(secondTimestamp)); + } finally { + if (second != null) { + second.close(); + } + } + + assertThat((Integer) oldestIteratorTimestampMetric.metricValue(), nullValue()); + } + 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 994a410bc5ef7..68ec8f1c79832 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 @@ -654,6 +654,42 @@ public void shouldTimeIteratorDuration() { assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(3.0 * TimeUnit.MILLISECONDS.toNanos(1))); } + @Test + public void shouldTrackOldestOpenIteratorTimestamp() { + when(innerStore.backwardFetch(KEY_BYTES)).thenReturn(KeyValueIterators.emptyIterator()); + init(); + + final KafkaMetric oldestIteratorTimestampMetric = metric("oldest-iterator-open-since-ms"); + assertThat(oldestIteratorTimestampMetric, not(nullValue())); + + assertThat(oldestIteratorTimestampMetric.metricValue(), nullValue()); + + KeyValueIterator, String> second = null; + final long secondTimestamp; + try { + try (final KeyValueIterator, String> first = store.backwardFetch(KEY)) { + final long oldestTimestamp = mockTime.milliseconds(); + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(oldestTimestamp)); + mockTime.sleep(100); + + // open a second iterator before closing the first to test that we still produce the first iterator's timestamp + second = store.backwardFetch(KEY); + secondTimestamp = mockTime.milliseconds(); + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(oldestTimestamp)); + mockTime.sleep(100); + } + + // now that the first iterator is closed, check that the timestamp has advanced to the still open second iterator + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(secondTimestamp)); + } finally { + if (second != null) { + second.close(); + } + } + + assertThat((Integer) oldestIteratorTimestampMetric.metricValue(), nullValue()); + } + 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 9fbd6ea93caf7..32d88e8c8ce0d 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 @@ -487,4 +487,40 @@ public void shouldTimeIteratorDuration() { assertThat((double) iteratorDurationAvgMetric.metricValue(), equalTo(2.5 * TimeUnit.MILLISECONDS.toNanos(1))); assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(3.0 * TimeUnit.MILLISECONDS.toNanos(1))); } + + @Test + public void shouldTrackOldestOpenIteratorTimestamp() { + when(inner.all()).thenReturn(KeyValueIterators.emptyIterator()); + init(); + + final KafkaMetric oldestIteratorTimestampMetric = metric("oldest-iterator-open-since-ms"); + assertThat(oldestIteratorTimestampMetric, not(nullValue())); + + assertThat(oldestIteratorTimestampMetric.metricValue(), nullValue()); + + KeyValueIterator> second = null; + final long secondTimestamp; + try { + try (final KeyValueIterator> first = metered.all()) { + final long oldestTimestamp = mockTime.milliseconds(); + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(oldestTimestamp)); + mockTime.sleep(100); + + // open a second iterator before closing the first to test that we still produce the first iterator's timestamp + second = metered.all(); + secondTimestamp = mockTime.milliseconds(); + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(oldestTimestamp)); + mockTime.sleep(100); + } + + // now that the first iterator is closed, check that the timestamp has advanced to the still open second iterator + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(secondTimestamp)); + } finally { + if (second != null) { + second.close(); + } + } + + assertThat((Integer) oldestIteratorTimestampMetric.metricValue(), nullValue()); + } } 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 09159427c5ac3..1e2a765bd021e 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 @@ -426,6 +426,48 @@ public void shouldTimeIteratorDuration() { assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(3.0 * TimeUnit.MILLISECONDS.toNanos(1))); } + @Test + public void shouldTrackOldestOpenIteratorTimestamp() { + 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 oldestIteratorTimestampMetric = getMetric("oldest-iterator-open-since-ms"); + assertThat(oldestIteratorTimestampMetric, not(nullValue())); + + assertThat(oldestIteratorTimestampMetric.metricValue(), nullValue()); + + final QueryResult> first = store.query(query, bound, config); + VersionedRecordIterator secondIterator = null; + final long secondTime; + try { + try (final VersionedRecordIterator iterator = first.getResult()) { + final long oldestTimestamp = mockTime.milliseconds(); + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(oldestTimestamp)); + mockTime.sleep(100); + + // open a second iterator before closing the first to test that we still produce the first iterator's timestamp + final QueryResult> second = store.query(query, bound, config); + secondIterator = second.getResult(); + secondTime = mockTime.milliseconds(); + + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(oldestTimestamp)); + mockTime.sleep(100); + } + + // now that the first iterator is closed, check that the timestamp has advanced to the still open second iterator + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(secondTime)); + } finally { + if (secondIterator != null) { + secondIterator.close(); + } + } + + assertThat((Integer) oldestIteratorTimestampMetric.metricValue(), nullValue()); + } + 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 c57bd1956f4d3..46e7d6f69cd54 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 @@ -495,6 +495,42 @@ public void shouldTimeIteratorDuration() { assertThat((double) iteratorDurationMaxMetric.metricValue(), equalTo(3.0 * TimeUnit.MILLISECONDS.toNanos(1))); } + @Test + public void shouldTrackOldestOpenIteratorTimestamp() { + when(innerStoreMock.all()).thenReturn(KeyValueIterators.emptyIterator()); + store.init((StateStoreContext) context, store); + + final KafkaMetric oldestIteratorTimestampMetric = metric("oldest-iterator-open-since-ms"); + assertThat(oldestIteratorTimestampMetric, not(nullValue())); + + assertThat(oldestIteratorTimestampMetric.metricValue(), nullValue()); + + KeyValueIterator, String> second = null; + final long secondTimestamp; + try { + try (final KeyValueIterator, String> first = store.all()) { + final long oldestTimestamp = mockTime.milliseconds(); + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(oldestTimestamp)); + mockTime.sleep(100); + + // open a second iterator before closing the first to test that we still produce the first iterator's timestamp + second = store.all(); + secondTimestamp = mockTime.milliseconds(); + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(oldestTimestamp)); + mockTime.sleep(100); + } + + // now that the first iterator is closed, check that the timestamp has advanced to the still open second iterator + assertThat((Long) oldestIteratorTimestampMetric.metricValue(), equalTo(secondTimestamp)); + } finally { + if (second != null) { + second.close(); + } + } + + assertThat((Integer) oldestIteratorTimestampMetric.metricValue(), nullValue()); + } + private KafkaMetric metric(final String name) { return metrics.metric(new MetricName(name, STORE_LEVEL_GROUP, "", tags)); }