Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Function;

import static org.apache.kafka.common.utils.Utils.mkEntry;
Expand Down Expand Up @@ -93,6 +94,8 @@ public class MeteredKeyValueStore<K, V>
private StreamsMetricsImpl streamsMetrics;
private TaskId taskId;

protected AtomicInteger numOpenIterators = new AtomicInteger(0);

@SuppressWarnings("rawtypes")
private final Map<Class, QueryHandler> queryHandlers =
mkMap(
Expand Down Expand Up @@ -162,6 +165,8 @@ 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);
StateStoreMetrics.addNumOpenIteratorsGauge(taskId.toString(), metricsScope, name(), streamsMetrics,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would it be better to add a Sensor that allows us to track the different metrics in one go?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can Sensors track Gauges? I don't believe that they can.

(config, now) -> numOpenIterators.get());
}

protected Serde<V> prepareValueSerdeForStore(final Serde<V> valueSerde, final SerdeGetter getter) {
Expand Down Expand Up @@ -460,6 +465,7 @@ private MeteredKeyValueIterator(final KeyValueIterator<Bytes, byte[]> iter,
this.iter = iter;
this.sensor = sensor;
this.startNs = time.nanoseconds();
numOpenIterators.incrementAndGet();
}

@Override
Expand All @@ -481,6 +487,7 @@ public void close() {
iter.close();
} finally {
sensor.record(time.nanoseconds() - startNs);
numOpenIterators.decrementAndGet();
}
}

Expand All @@ -504,6 +511,7 @@ private MeteredKeyValueTimestampedIterator(final KeyValueIterator<Bytes, byte[]>
this.sensor = sensor;
this.valueDeserializer = valueDeserializer;
this.startNs = time.nanoseconds();
numOpenIterators.incrementAndGet();
}

@Override
Expand All @@ -525,6 +533,7 @@ public void close() {
iter.close();
} finally {
sensor.record(time.nanoseconds() - startNs);
numOpenIterators.decrementAndGet();
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
*/
package org.apache.kafka.streams.state.internals;

import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Function;
import org.apache.kafka.streams.state.VersionedRecordIterator;
import org.apache.kafka.streams.state.VersionedRecord;
Expand All @@ -24,18 +25,25 @@ public class MeteredMultiVersionedKeyQueryIterator<V> implements VersionedRecord

private final VersionedRecordIterator<byte[]> iterator;
private final Function<VersionedRecord<byte[]>, VersionedRecord<V>> deserializeValue;

private final AtomicInteger numOpenIterators;

public MeteredMultiVersionedKeyQueryIterator(final VersionedRecordIterator<byte[]> iterator,
final Function<VersionedRecord<byte[]>, VersionedRecord<V>> deserializeValue) {
final Function<VersionedRecord<byte[]>, VersionedRecord<V>> deserializeValue,
final AtomicInteger numOpenIterators) {
this.iterator = iterator;
this.deserializeValue = deserializeValue;
this.numOpenIterators = numOpenIterators;
numOpenIterators.incrementAndGet();
}


@Override
public void close() {
iterator.close();
try {
iterator.close();
} finally {
numOpenIterators.decrementAndGet();
}
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@

import java.util.Map;
import java.util.Objects;
import java.util.concurrent.atomic.AtomicInteger;

import static org.apache.kafka.common.utils.Utils.mkEntry;
import static org.apache.kafka.common.utils.Utils.mkMap;
Expand All @@ -70,6 +71,8 @@ public class MeteredSessionStore<K, V>
private InternalProcessorContext<?, ?> context;
private TaskId taskId;

private AtomicInteger numOpenIterators = new AtomicInteger(0);

@SuppressWarnings("rawtypes")
private final Map<Class, QueryHandler> queryHandlers =
mkMap(
Expand Down Expand Up @@ -131,6 +134,8 @@ 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);
StateStoreMetrics.addNumOpenIteratorsGauge(taskId.toString(), metricsScope, name(), streamsMetrics,
(config, now) -> numOpenIterators.get());
}


Expand Down Expand Up @@ -248,7 +253,8 @@ public KeyValueIterator<Windowed<K>, V> fetch(final K key) {
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time);
time,
numOpenIterators);
}

@Override
Expand All @@ -260,7 +266,8 @@ public KeyValueIterator<Windowed<K>, V> backwardFetch(final K key) {
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time
time,
numOpenIterators
);
}

Expand All @@ -273,7 +280,8 @@ public KeyValueIterator<Windowed<K>, V> fetch(final K keyFrom,
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time);
time,
numOpenIterators);
}

@Override
Expand All @@ -285,7 +293,8 @@ public KeyValueIterator<Windowed<K>, V> backwardFetch(final K keyFrom,
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time
time,
numOpenIterators
);
}

Expand All @@ -304,7 +313,8 @@ public KeyValueIterator<Windowed<K>, V> findSessions(final K key,
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time);
time,
numOpenIterators);
}

@Override
Expand All @@ -323,7 +333,8 @@ public KeyValueIterator<Windowed<K>, V> backwardFindSessions(final K key,
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time
time,
numOpenIterators
);
}

Expand All @@ -344,7 +355,8 @@ public KeyValueIterator<Windowed<K>, V> findSessions(final K keyFrom,
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time);
time,
numOpenIterators);
}

@Override
Expand All @@ -356,7 +368,8 @@ public KeyValueIterator<Windowed<K>, V> findSessions(final long earliestSessionE
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time);
time,
numOpenIterators);
}

@Override
Expand All @@ -377,7 +390,8 @@ public KeyValueIterator<Windowed<K>, V> backwardFindSessions(final K keyFrom,
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time
time,
numOpenIterators
);
}

Expand Down Expand Up @@ -447,7 +461,8 @@ private <R> QueryResult<R> runRangeQuery(final Query<R> query,
streamsMetrics,
serdes::keyFrom,
StoreQueryUtils.getDeserializeValue(serdes, wrapped()),
time
time,
numOpenIterators
);
final QueryResult<MeteredWindowedKeyValueIterator<K, V>> typedQueryResult =
InternalQueryResultUtil.copyAndSubstituteDeserializedResult(rawResult, typedResult);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -325,6 +325,7 @@ private MeteredTimestampedKeyValueStoreIterator(final KeyValueIterator<Bytes, by
this.valueAndTimestampDeserializer = valueAndTimestampDeserializer;
this.startNs = time.nanoseconds();
this.returnPlainValue = returnPlainValue;
numOpenIterators.incrementAndGet();
}

@Override
Expand All @@ -350,6 +351,7 @@ public void close() {
iter.close();
} finally {
sensor.record(time.nanoseconds() - startNs);
numOpenIterators.decrementAndGet();
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -264,7 +264,7 @@ private <R> QueryResult<R> runMultiVersionedKeyQuery(final Query<R> query, final
final QueryResult<VersionedRecordIterator<byte[]>> rawResult = wrapped().query(rawKeyQuery, positionBound, config);
if (rawResult.isSuccess()) {
final MeteredMultiVersionedKeyQueryIterator<V> typedResult =
new MeteredMultiVersionedKeyQueryIterator<V>(rawResult.getResult(), StoreQueryUtils.getDeserializeValue(plainValueSerdes));
new MeteredMultiVersionedKeyQueryIterator<V>(rawResult.getResult(), StoreQueryUtils.getDeserializeValue(plainValueSerdes), numOpenIterators);
final QueryResult<MeteredMultiVersionedKeyQueryIterator<V>> typedQueryResult =
InternalQueryResultUtil.copyAndSubstituteDeserializedResult(rawResult, typedResult);
result = (QueryResult<R>) typedQueryResult;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@

import java.util.Map;
import java.util.Objects;
import java.util.concurrent.atomic.AtomicInteger;

import static org.apache.kafka.common.utils.Utils.mkEntry;
import static org.apache.kafka.common.utils.Utils.mkMap;
Expand All @@ -74,6 +75,8 @@ public class MeteredWindowStore<K, V>
private InternalProcessorContext<?, ?> context;
private TaskId taskId;

private AtomicInteger numOpenIterators = new AtomicInteger(0);

@SuppressWarnings("rawtypes")
private final Map<Class, QueryHandler> queryHandlers =
mkMap(
Expand Down Expand Up @@ -150,6 +153,8 @@ 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);
StateStoreMetrics.addNumOpenIteratorsGauge(taskId.toString(), metricsScope, name(), streamsMetrics,
(config, now) -> numOpenIterators.get());
}

@Deprecated
Expand Down Expand Up @@ -236,7 +241,8 @@ public WindowStoreIterator<V> fetch(final K key,
fetchSensor,
streamsMetrics,
serdes::valueFrom,
time
time,
numOpenIterators
);
}

Expand All @@ -250,7 +256,8 @@ public WindowStoreIterator<V> backwardFetch(final K key,
fetchSensor,
streamsMetrics,
serdes::valueFrom,
time
time,
numOpenIterators
);
}

Expand All @@ -269,7 +276,8 @@ public KeyValueIterator<Windowed<K>, V> fetch(final K keyFrom,
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time);
time,
numOpenIterators);
}

@Override
Expand All @@ -287,7 +295,8 @@ public KeyValueIterator<Windowed<K>, V> backwardFetch(final K keyFrom,
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time);
time,
numOpenIterators);
}

@Override
Expand All @@ -299,7 +308,8 @@ public KeyValueIterator<Windowed<K>, V> fetchAll(final long timeFrom,
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time);
time,
numOpenIterators);
}

@Override
Expand All @@ -311,7 +321,8 @@ public KeyValueIterator<Windowed<K>, V> backwardFetchAll(final long timeFrom,
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time);
time,
numOpenIterators);
}

@Override
Expand All @@ -322,7 +333,8 @@ public KeyValueIterator<Windowed<K>, V> all() {
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time
time,
numOpenIterators
);
}

Expand All @@ -334,7 +346,8 @@ public KeyValueIterator<Windowed<K>, V> backwardAll() {
streamsMetrics,
serdes::keyFrom,
serdes::valueFrom,
time
time,
numOpenIterators
);
}

Expand Down Expand Up @@ -410,7 +423,8 @@ private <R> QueryResult<R> runRangeQuery(final Query<R> query,
streamsMetrics,
serdes::keyFrom,
getDeserializeValue(serdes, wrapped()),
time
time,
numOpenIterators
);
final QueryResult<MeteredWindowedKeyValueIterator<K, V>> typedQueryResult =
InternalQueryResultUtil.copyAndSubstituteDeserializedResult(rawResult, typedResult);
Expand Down Expand Up @@ -459,7 +473,8 @@ private <R> QueryResult<R> runKeyQuery(final Query<R> query,
fetchSensor,
streamsMetrics,
getDeserializeValue(serdes, wrapped()),
time
time,
numOpenIterators
);
final QueryResult<MeteredWindowStoreIterator<V>> typedQueryResult =
InternalQueryResultUtil.copyAndSubstituteDeserializedResult(rawResult, typedResult);
Expand Down
Loading