aliehsaeedii commented on code in PR #22975:
URL: https://github.com/apache/kafka/pull/22975#discussion_r3726259572
##########
streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStoreWithHeadersTest.java:
##########
@@ -619,6 +622,63 @@ public void
shouldDecrementOpenIteratorsTwiceWhenClosedTwiceForTimestampedWindow
assertEquals(-1L, (Long) openIterators.metricValue());
}
+ // The window store previously had no iterator-duration coverage at all.
This mirrors the
+ // session/KV shouldTimeIteratorDuration: it goes through store.all() ->
the KeyValueIterator
+ // sibling, whose close() records the operation (fetch) and
iterator-duration sensors via the
+ // shared AbstractMeteredIterator lifecycle.
+ @Test
+ public void shouldTimeIteratorDuration() {
+ setUp();
+ store.init(context, store);
+ when(innerStoreMock.all()).thenReturn(windowRangeIterator(List.of()));
+
+ final KafkaMetric iteratorDurationAvg =
metric("iterator-duration-avg");
+ final KafkaMetric iteratorDurationMax =
metric("iterator-duration-max");
+ assertEquals(Double.NaN, (Double) iteratorDurationAvg.metricValue());
+ assertEquals(Double.NaN, (Double) iteratorDurationMax.metricValue());
+
+ try (KeyValueIterator<Windowed<String>, ValueTimestampHeaders<String>>
iterator = store.all()) {
+ // nothing to iterate; just hold it open, then close
+ mockTime.sleep(2);
+ }
+
+ assertEquals(2.0 * TimeUnit.MILLISECONDS.toNanos(1), (double)
iteratorDurationAvg.metricValue());
+ assertEquals(2.0 * TimeUnit.MILLISECONDS.toNanos(1), (double)
iteratorDurationMax.metricValue());
Review Comment:
One sample makes avg and max identical, so this doesn't really pin avg. The
KV sibling `shouldTimeIteratorDuration` opens two iterators (2ms then 3ms) so
avg is 2.5ms and max is 3ms. Should we do mirroring here.
##########
streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractMeteredIterator.java:
##########
@@ -0,0 +1,97 @@
+/*
+ * 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;
+
+import org.apache.kafka.common.metrics.Sensor;
+import org.apache.kafka.common.utils.Time;
+import org.apache.kafka.streams.state.KeyValueIterator;
+
+import java.util.Set;
+import java.util.concurrent.atomic.LongAdder;
+
+/**
+ * Shared metering lifecycle for the metered iterators of the {@code
Metered*WithHeaders} stores,
+ * whatever result type they yield: the {@code KeyValueIterator}s returned by
the store's own range/
+ * fetch/find methods and the {@code ReadOnlyRecordIterator}s that back the
headers-aware IQv2
+ * range/window/session query types.
+ *
+ * <p>Every such iterator opens over a raw {@code KeyValueIterator<RawKey,
byte[]>} and needs the
+ * same bookkeeping: stamp the open time (for the {@code
oldest-iterator-open-since-ms} metric),
+ * register in {@code numOpenIterators}/{@code openIterators}, and on {@link
#close()} record the
+ * operation and iterator-duration sensors and deregister. This base is
deliberately result-type
+ * agnostic -- it implements only {@link MeteredIterator} and does not bind
the yielded key/value
+ * types -- so each subclass declares its own result interface (a {@code
KeyValueIterator} or a
+ * {@code ReadOnlyRecordIterator}) and implements just the parts that
genuinely differ: the
+ * deserializing {@code next()} (and, for the {@code KeyValueIterator}s, a
peeking {@code hasNext()}
+ * and {@code peekNextKey()}).
+ *
+ * @param <RawKey> the raw iterator's key type
+ */
+abstract class AbstractMeteredIterator<RawKey> implements MeteredIterator {
Review Comment:
`MeteredWindowedKeyValueIterator`, `MeteredWindowStoreIterator` and
`MeteredKeyValueStoreIterator` still hand-roll this exact lifecycle, field for
field. The first is the base of `MeteredWindowedKeyValueWithHeadersIterator`,
so one `Metered*WithHeaders` iterator is still left out. Should we do a
follow-up making them extend this class, after which the javadoc's
`Metered*WithHeaders` scoping can go.
##########
streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractMeteredIterator.java:
##########
@@ -0,0 +1,97 @@
+/*
+ * 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;
+
+import org.apache.kafka.common.metrics.Sensor;
+import org.apache.kafka.common.utils.Time;
+import org.apache.kafka.streams.state.KeyValueIterator;
+
+import java.util.Set;
+import java.util.concurrent.atomic.LongAdder;
+
+/**
+ * Shared metering lifecycle for the metered iterators of the {@code
Metered*WithHeaders} stores,
+ * whatever result type they yield: the {@code KeyValueIterator}s returned by
the store's own range/
+ * fetch/find methods and the {@code ReadOnlyRecordIterator}s that back the
headers-aware IQv2
+ * range/window/session query types.
+ *
+ * <p>Every such iterator opens over a raw {@code KeyValueIterator<RawKey,
byte[]>} and needs the
+ * same bookkeeping: stamp the open time (for the {@code
oldest-iterator-open-since-ms} metric),
+ * register in {@code numOpenIterators}/{@code openIterators}, and on {@link
#close()} record the
+ * operation and iterator-duration sensors and deregister. This base is
deliberately result-type
+ * agnostic -- it implements only {@link MeteredIterator} and does not bind
the yielded key/value
+ * types -- so each subclass declares its own result interface (a {@code
KeyValueIterator} or a
+ * {@code ReadOnlyRecordIterator}) and implements just the parts that
genuinely differ: the
+ * deserializing {@code next()} (and, for the {@code KeyValueIterator}s, a
peeking {@code hasNext()}
+ * and {@code peekNextKey()}).
+ *
+ * @param <RawKey> the raw iterator's key type
+ */
+abstract class AbstractMeteredIterator<RawKey> implements MeteredIterator {
+
+ final KeyValueIterator<RawKey, byte[]> iter;
+ private final Sensor operationSensor;
+ private final Sensor iteratorSensor;
+ private final Time time;
+ private final LongAdder numOpenIterators;
+ private final Set<MeteredIterator> openIterators;
+ private final long startNs;
+ private final long startTimestampMs;
+
+ AbstractMeteredIterator(final KeyValueIterator<RawKey, byte[]> iter,
+ final Sensor operationSensor,
+ final Sensor iteratorSensor,
+ final Time time,
+ final LongAdder numOpenIterators,
+ final Set<MeteredIterator> openIterators) {
+ this.iter = iter;
+ this.operationSensor = operationSensor;
+ this.iteratorSensor = iteratorSensor;
+ this.time = time;
+ this.numOpenIterators = numOpenIterators;
+ this.openIterators = openIterators;
+ this.startNs = time.nanoseconds();
+ this.startTimestampMs = time.milliseconds();
+ numOpenIterators.increment();
+ openIterators.add(this);
+ }
+
+ @Override
+ public long startTimestamp() {
Review Comment:
`openIterators.add(this)` in the constructor calls this method through the
set's comparator, i.e. on a half-built object. Worth making it `final` so a
subclass can't override it with something that reads its own not-yet-assigned
state.
##########
streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreWithHeadersTest.java:
##########
@@ -667,6 +667,37 @@ public void shouldTimeIteratorDuration() {
assertTrue((Double) iteratorDurationMetric.metricValue() > 0.0);
}
+ // The above shouldTimeIteratorDuration goes through store.fetch() -> the
KeyValueIterator sibling.
+ // This pins the same close()-path recording for the
ReadOnlyRecordIterator that backs
+ // TimestampedWindowRangeWithHeadersQuery.withKey, whose close() records
both the operation sensor
+ // (fetch) and the iterator-duration sensor via the shared
AbstractMeteredIterator lifecycle.
+ @SuppressWarnings({"unchecked", "rawtypes"})
+ @Test
+ public void
shouldTimeIteratorDurationForTimestampedWindowRangeWithHeadersQuery() {
+ setUp();
+ init();
+
+ when(innerStore.query(any(), any(PositionBound.class),
any(QueryConfig.class)))
+ .thenReturn((QueryResult) QueryResult.forResult(new
KeyValueIteratorStub<>(
+ Collections.<KeyValue<Windowed<Bytes>,
byte[]>>emptyList().iterator())));
+
+ final KafkaMetric iteratorDurationMetric =
metric("iterator-duration-avg");
+ final KafkaMetric fetchLatencyMetric = metric("fetch-latency-avg");
+
+ final QueryResult<ReadOnlyRecordIterator<Windowed<String>, String>>
result = store.query(
+ TimestampedWindowRangeWithHeadersQuery.<String,
String>withKey(KEY),
+ PositionBound.unbounded(),
+ new QueryConfig(false));
+ assertTrue(result.isSuccess());
+ try (ReadOnlyRecordIterator<Windowed<String>, String> iterator =
result.getResult()) {
+ // nothing to iterate; just hold it open, then close
+ mockTime.sleep(100L);
+ }
+
+ assertTrue((Double) iteratorDurationMetric.metricValue() > 0.0);
Review Comment:
`mockTime` makes this deterministic, so assert the exact `100.0 *
TimeUnit.MILLISECONDS.toNanos(1)` like the two sibling tests this PR adds,
instead of `> 0.0`.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]