diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index 8a74ae6df248c..3fc7af0f0aab6 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -2229,7 +2229,7 @@ private void recordRatio(final long now, final WindowedSum windowedSum, final Se if (runOnceLatencyWindow > 0.0) { final double latencyWindow = windowedSum.measure(metricsConfig, now); - ratioSensor.record(latencyWindow / runOnceLatencyWindow); + ratioSensor.record(latencyWindow / runOnceLatencyWindow, now); } else { ratioSensor.record(0.0, now); } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/TaskMetrics.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/TaskMetrics.java index 1912db0f481e4..7e85f7c6a85df 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/TaskMetrics.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/TaskMetrics.java @@ -28,7 +28,7 @@ import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.TOTAL_SUFFIX; import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.addAvgAndMaxToSensor; import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.addInvocationRateAndCountToSensor; -import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.addInvocationRateToSensor; +import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.addRateOfSumAndSumMetricsToSensor; import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.addSumMetricToSensor; import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.addValueMetricToSensor; @@ -189,7 +189,7 @@ public static Sensor restoreSensor(final String threadId, final String taskId, final StreamsMetricsImpl streamsMetrics, final Sensor... parentSensor) { - return invocationRateAndTotalSensor( + return rateAndTotalSensor( threadId, taskId, RESTORE, @@ -205,7 +205,7 @@ public static Sensor updateSensor(final String threadId, final String taskId, final StreamsMetricsImpl streamsMetrics, final Sensor... parentSensor) { - return invocationRateAndTotalSensor( + return rateAndTotalSensor( threadId, taskId, UPDATE, @@ -250,7 +250,7 @@ public static Sensor recordLatenessSensor(final String threadId, public static Sensor droppedRecordsSensor(final String threadId, final String taskId, final StreamsMetricsImpl streamsMetrics) { - return invocationRateAndTotalSensor( + return rateAndTotalSensor( threadId, taskId, DROPPED_RECORDS, @@ -281,19 +281,20 @@ private static Sensor invocationRateAndCountSensor(final String threadId, return sensor; } - private static Sensor invocationRateAndTotalSensor(final String threadId, - final String taskId, - final String operation, - final String descriptionOfRate, - final String descriptionOfTotal, - final RecordingLevel recordingLevel, - final StreamsMetricsImpl streamsMetrics, - final Sensor... parentSensors) { + private static Sensor rateAndTotalSensor( + final String threadId, + final String taskId, + final String operation, + final String descriptionOfRate, + final String descriptionOfTotal, + final RecordingLevel recordingLevel, + final StreamsMetricsImpl streamsMetrics, + final Sensor... parentSensors + ) { final Sensor sensor = streamsMetrics.taskLevelSensor(threadId, taskId, operation, recordingLevel, parentSensors); final Map tags = streamsMetrics.taskLevelTagMap(threadId, taskId); - addInvocationRateToSensor(sensor, TASK_LEVEL_GROUP, tags, operation, descriptionOfRate); - addSumMetricToSensor(sensor, TASK_LEVEL_GROUP, tags, operation, true, descriptionOfTotal); + addRateOfSumAndSumMetricsToSensor(sensor, TASK_LEVEL_GROUP, tags, operation, descriptionOfRate, descriptionOfTotal); return sensor; } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractDualSchemaRocksDBSegmentedBytesStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractDualSchemaRocksDBSegmentedBytesStore.java index a90d87b907744..963fd8c6b8674 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractDualSchemaRocksDBSegmentedBytesStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractDualSchemaRocksDBSegmentedBytesStore.java @@ -189,7 +189,7 @@ public void put(final Bytes rawBaseKey, final S segment = segments.getOrCreateSegmentIfLive(segmentId, internalProcessorContext, observedStreamTime); if (segment == null) { - expiredRecordSensor.record(1.0d, internalProcessorContext.currentSystemTimeMs()); + expiredRecordSensor.record(); LOG.warn("Skipping record for expired segment."); } else { synchronized (position) { diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractRocksDBSegmentedBytesStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractRocksDBSegmentedBytesStore.java index 0d3c25783ae55..960c849bc6909 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractRocksDBSegmentedBytesStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractRocksDBSegmentedBytesStore.java @@ -254,7 +254,7 @@ public void put(final Bytes key, final long segmentId = segments.segmentId(timestamp); final S segment = segments.getOrCreateSegmentIfLive(segmentId, internalProcessorContext, observedStreamTime); if (segment == null) { - expiredRecordSensor.record(1.0d, internalProcessorContext.currentSystemTimeMs()); + expiredRecordSensor.record(); } else { synchronized (position) { // Transactional puts stage their position too, so READ_COMMITTED never sees a position ahead of the data. diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemorySessionStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemorySessionStore.java index fa81a903dfa4d..caa0a6ddeb898 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemorySessionStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemorySessionStore.java @@ -159,8 +159,8 @@ public void init(final StateStoreContext stateStoreContext, } removeExpiredSegments(); if (expiredRecords > 0) { - if (expiredRecordSensor != null && context != null) { - expiredRecordSensor.record(expiredRecords, context.currentSystemTimeMs()); + if (expiredRecordSensor != null) { + expiredRecordSensor.record(expiredRecords); } LOG.warn("Skipping {} records for expired segments.", expiredRecords); } @@ -205,8 +205,8 @@ public void put(final Windowed sessionKey, final byte[] aggregate) { if (windowEndTimestamp <= observedStreamTime - retentionPeriod) { // The provided context is not required to implement InternalProcessorContext, // If it doesn't, we can't record this metric (in fact, we wouldn't have even initialized it). - if (expiredRecordSensor != null && context != null) { - expiredRecordSensor.record(1.0d, context.currentSystemTimeMs()); + if (expiredRecordSensor != null) { + expiredRecordSensor.record(); } LOG.warn("Skipping record for expired segment."); } else if (transactionBuffer != null) { diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryWindowStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryWindowStore.java index aa61bbc0a48bd..cb6e26445e710 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryWindowStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryWindowStore.java @@ -155,7 +155,7 @@ public void init(final StateStoreContext stateStoreContext, } removeExpiredSegments(); if (expiredRecords > 0) { - expiredRecordSensor.record(expiredRecords, internalProcessorContext.currentSystemTimeMs()); + expiredRecordSensor.record(expiredRecords); LOG.warn("Skipping {} records for expired segments.", expiredRecords); } } @@ -195,7 +195,7 @@ public void put(final Bytes key, final byte[] value, final long windowStartTimes synchronized (position) { if (windowStartTimestamp <= observedStreamTime - retentionPeriod) { - expiredRecordSensor.record(1.0d, internalProcessorContext.currentSystemTimeMs()); + expiredRecordSensor.record(); LOG.warn("Skipping record for expired segment."); } else if (transactionBuffer != null) { if (value != null) { diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBVersionedStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBVersionedStore.java index 6477427e66ffe..82ee65d7839aa 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBVersionedStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBVersionedStore.java @@ -132,7 +132,7 @@ public long put(final Bytes key, final byte[] value, final long timestamp) { synchronized (position) { if (timestamp < observedStreamTime - gracePeriod) { - expiredRecordSensor.record(1.0d, internalProcessorContext.currentSystemTimeMs()); + expiredRecordSensor.record(); LOG.warn("Skipping record for expired put."); StoreQueryUtils.updatePosition(position, internalProcessorContext); return PUT_RETURN_CODE_NOT_PUT; @@ -160,7 +160,7 @@ public VersionedRecord delete(final Bytes key, final long timestamp) { synchronized (position) { if (timestamp < observedStreamTime - gracePeriod) { - expiredRecordSensor.record(1.0d, internalProcessorContext.currentSystemTimeMs()); + expiredRecordSensor.record(); LOG.warn("Skipping record for expired delete."); return null; } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StandbyTaskTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StandbyTaskTest.java index 84d1f218cfc33..3ddf2fb8e5cab 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StandbyTaskTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StandbyTaskTest.java @@ -75,7 +75,6 @@ import static org.hamcrest.Matchers.empty; import static org.hamcrest.Matchers.is; import static org.hamcrest.Matchers.isA; -import static org.hamcrest.Matchers.not; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.mockito.ArgumentMatchers.any; @@ -487,12 +486,26 @@ public void shouldRecordRestoredRecords() { task.recordRestoration(time, 25L, false); assertThat(totalMetric.metricValue(), equalTo(25.0)); - assertThat(rateMetric.metricValue(), not(0.0)); + // the rate measures updated records per second, not update batches per second; with no time + // elapsed the rate window is (metrics.num.samples - 1) * metrics.sample.window.ms == 30s + assertEquals( + 25.0 / 30.0, + ((Number) rateMetric.metricValue()).doubleValue(), + 0.0001d, + "update-rate must measure updated records per second, not update batches per second; " + + "counting batches would give 1/30 == 0.03333 (KAFKA-20877)" + ); task.recordRestoration(time, 50L, false); assertThat(totalMetric.metricValue(), equalTo(75.0)); - assertThat(rateMetric.metricValue(), not(0.0)); + assertEquals( + 75.0 / 30.0, + ((Number) rateMetric.metricValue()).doubleValue(), + 0.0001d, + "update-rate must measure updated records per second, not update batches per second; " + + "counting batches would give 2/30 == 0.06666 (KAFKA-20877)" + ); } private KafkaMetric getMetric(final String operation, diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java index 48639fad16ac2..85fc510ae8d26 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java @@ -893,13 +893,27 @@ public void shouldRecordRestoredRecords() { task.recordRestoration(time, 25L, false); assertThat(totalMetric.metricValue(), equalTo(25.0)); - assertThat(rateMetric.metricValue(), not(0.0)); + // the rate measures restored records per second, not restore batches per second; with no time + // elapsed the rate window is (metrics.num.samples - 1) * metrics.sample.window.ms == 30s + assertEquals( + 25.0 / 30.0, + ((Number) rateMetric.metricValue()).doubleValue(), + 0.0001d, + "restore-rate must measure restored records per second, not restore batches per second; " + + "counting batches would give 1/30 == 0.03333 (KAFKA-20877)" + ); assertThat(remainMetric.metricValue(), equalTo(75.0)); task.recordRestoration(time, 50L, false); assertThat(totalMetric.metricValue(), equalTo(75.0)); - assertThat(rateMetric.metricValue(), not(0.0)); + assertEquals( + 75.0 / 30.0, + ((Number) rateMetric.metricValue()).doubleValue(), + 0.0001d, + "restore-rate must measure restored records per second, not restore batches per second; " + + "counting batches would give 2/30 == 0.06666 (KAFKA-20877)" + ); assertThat(remainMetric.metricValue(), equalTo(25.0)); } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/TaskMetricsTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/TaskMetricsTest.java index 13add9eabc806..cfdc7b7196eed 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/TaskMetricsTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/TaskMetricsTest.java @@ -240,21 +240,12 @@ public void shouldGetDroppedRecordsSensor() { try (final MockedStatic streamsMetricsStaticMock = mockStatic(StreamsMetricsImpl.class)) { final Sensor sensor = TaskMetrics.droppedRecordsSensor(THREAD_ID, TASK_ID, streamsMetrics); streamsMetricsStaticMock.verify( - () -> StreamsMetricsImpl.addInvocationRateToSensor( + () -> StreamsMetricsImpl.addRateOfSumAndSumMetricsToSensor( expectedSensor, TASK_LEVEL_GROUP, tagMap, operation, - rateDescription - ) - ); - streamsMetricsStaticMock.verify( - () -> StreamsMetricsImpl.addSumMetricToSensor( - expectedSensor, - TASK_LEVEL_GROUP, - tagMap, - operation, - true, + rateDescription, totalDescription ) ); diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractDualSchemaRocksDBSegmentedBytesStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractDualSchemaRocksDBSegmentedBytesStoreTest.java index 3e51e20a35d43..513eb01cfeae8 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractDualSchemaRocksDBSegmentedBytesStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractDualSchemaRocksDBSegmentedBytesStoreTest.java @@ -88,7 +88,6 @@ import static org.hamcrest.Matchers.hasEntry; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; -import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -1586,8 +1585,6 @@ public void shouldMeasureExpiredRecords() { TestUtils.tempDirectory(), new StreamsConfig(streamsConfig) ); - final Time time = Time.SYSTEM; - context.setSystemTimeMs(time.milliseconds()); bytesStore.init(context, bytesStore); // write a record to advance stream time, with a high enough timestamp @@ -1622,7 +1619,16 @@ public void shouldMeasureExpiredRecords() { ) )); assertEquals(1.0, dropTotal.metricValue()); - assertNotEquals(0.0, dropRate.metricValue()); + // exactly one record was dropped, over the rate's default un-elapsed sampling window of + // (metrics.num.samples - 1) * metrics.sample.window.ms == 30s. The delta is generous because the + // window also grows by however long the store work takes between recording and reading the + // metric; it still separates one dropped record from none (0.0) and from two (0.06666). + assertEquals( + 1.0 / 30.0, + ((Number) dropRate.metricValue()).doubleValue(), + 0.005d, + "dropped-records-rate should reflect the single dropped record over the ~30s sampling window" + ); bytesStore.close(); } diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractRocksDBSegmentedBytesStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractRocksDBSegmentedBytesStoreTest.java index acd6d219cc95f..35ffb27cf2b80 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractRocksDBSegmentedBytesStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractRocksDBSegmentedBytesStoreTest.java @@ -95,7 +95,6 @@ import static org.hamcrest.Matchers.hasEntry; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; -import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -962,8 +961,6 @@ public void shouldMeasureExpiredRecords(final SegmentedBytesStore.KeySchema sche TestUtils.tempDirectory(), new StreamsConfig(streamsConfig) ); - final Time time = Time.SYSTEM; - context.setSystemTimeMs(time.milliseconds()); bytesStore.init(context, bytesStore); // write a record to advance stream time, with a high enough timestamp @@ -998,7 +995,16 @@ public void shouldMeasureExpiredRecords(final SegmentedBytesStore.KeySchema sche ) )); assertEquals(1.0, dropTotal.metricValue()); - assertNotEquals(0.0, dropRate.metricValue()); + // exactly one record was dropped, over the rate's default un-elapsed sampling window of + // (metrics.num.samples - 1) * metrics.sample.window.ms == 30s. The delta is generous because the + // window also grows by however long the store work takes between recording and reading the + // metric; it still separates one dropped record from none (0.0) and from two (0.06666). + assertEquals( + 1.0 / 30.0, + ((Number) dropRate.metricValue()).doubleValue(), + 0.005d, + "dropped-records-rate should reflect the single dropped record over the ~30s sampling window" + ); bytesStore.close(); } diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractSessionBytesStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractSessionBytesStoreTest.java index c0a35a8e4257a..c4059027dbcfb 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractSessionBytesStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractSessionBytesStoreTest.java @@ -71,7 +71,6 @@ import static org.hamcrest.Matchers.is; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; -import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -894,9 +893,7 @@ public void shouldMeasureExpiredRecords() { new StreamsConfig(streamsConfig), recordCollector ); - final Time time = Time.SYSTEM; context.setTime(1L); - context.setSystemTimeMs(time.milliseconds()); sessionStore.init(context, sessionStore); // Advance stream time by inserting record with large enough timestamp that records with timestamp 0 are expired @@ -931,7 +928,16 @@ public void shouldMeasureExpiredRecords() { ) )); assertEquals(1.0, dropTotal.metricValue()); - assertNotEquals(0.0, dropRate.metricValue()); + // exactly one record was dropped, over the rate's default un-elapsed sampling window of + // (metrics.num.samples - 1) * metrics.sample.window.ms == 30s. The delta is generous because the + // window also grows by however long the store work takes between recording and reading the + // metric; it still separates one dropped record from none (0.0) and from two (0.06666). + assertEquals( + 1.0 / 30.0, + ((Number) dropRate.metricValue()).doubleValue(), + 0.005d, + "dropped-records-rate should reflect the single dropped record over the ~30s sampling window" + ); sessionStore.close(); } diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractWindowBytesStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractWindowBytesStoreTest.java index edcb73bad02da..4372c8c641cb5 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractWindowBytesStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractWindowBytesStoreTest.java @@ -25,7 +25,6 @@ import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.common.utils.Bytes; import org.apache.kafka.common.utils.LogCaptureAppender; -import org.apache.kafka.common.utils.Time; import org.apache.kafka.common.utils.internals.LogContext; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsConfig; @@ -67,7 +66,6 @@ import static org.hamcrest.MatcherAssert.assertThat; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; -import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -1007,8 +1005,6 @@ public void shouldMeasureExpiredRecords() { new StreamsConfig(streamsConfig), recordCollector ); - final Time time = Time.SYSTEM; - context.setSystemTimeMs(time.milliseconds()); context.setTime(1L); windowStore.init(context, windowStore); @@ -1044,7 +1040,16 @@ public void shouldMeasureExpiredRecords() { ) )); assertEquals(1.0, dropTotal.metricValue()); - assertNotEquals(0.0, dropRate.metricValue()); + // exactly one record was dropped, over the rate's default un-elapsed sampling window of + // (metrics.num.samples - 1) * metrics.sample.window.ms == 30s. The delta is generous because the + // window also grows by however long the store work takes between recording and reading the + // metric; it still separates one dropped record from none (0.0) and from two (0.06666). + assertEquals( + 1.0 / 30.0, + ((Number) dropRate.metricValue()).doubleValue(), + 0.005d, + "dropped-records-rate should reflect the single dropped record over the ~30s sampling window" + ); windowStore.close(); } diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/InMemorySessionStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/InMemorySessionStoreTest.java index 305431899771f..024da3fedee5c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/InMemorySessionStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/InMemorySessionStoreTest.java @@ -17,9 +17,12 @@ package org.apache.kafka.streams.state.internals; import org.apache.kafka.common.IsolationLevel; +import org.apache.kafka.common.Metric; +import org.apache.kafka.common.MetricName; import org.apache.kafka.common.header.internals.RecordHeaders; import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.common.utils.Bytes; +import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.kstream.Windowed; import org.apache.kafka.streams.kstream.internals.SessionWindow; @@ -37,6 +40,8 @@ import org.junit.jupiter.api.Test; +import java.util.LinkedList; +import java.util.List; import java.util.Map; import java.util.Properties; import java.util.Set; @@ -318,4 +323,44 @@ private InMemorySessionStore openTransactionalSessionStore() { return store; } + + @Test + public void shouldMeasureExpiredRecordsDroppedDuringRestoreAsRecords() { + // The restore path reports every record skipped for an expired segment in a single sensor + // recording, so the rate has to reflect the number of records dropped rather than the number of + // recordings. Mirrors the same coverage for InMemoryWindowStore. + + final List> batch = new LinkedList<>(); + // advances observed stream time far enough that every record after it falls outside retention + batch.add(new KeyValue<>( + SessionKeySchema.toBinary(Bytes.wrap("on-time".getBytes()), 0L, 4 * RETENTION_PERIOD).get(), + Serdes.Long().serializer().serialize("", 0L))); + for (int key = 1; key <= 3; key++) { + batch.add(new KeyValue<>( + SessionKeySchema.toBinary(Bytes.wrap(("expired-" + key).getBytes()), 0L, 0L).get(), + Serdes.Long().serializer().serialize("", (long) key))); + } + + context.restore(sessionStore.name(), batch); + + final Map metrics = context.metrics().metrics(); + final Map tags = mkMap( + mkEntry("thread-id", Thread.currentThread().getName()), + mkEntry("task-id", "0_0") + ); + final Metric dropTotal = metrics.get( + new MetricName("dropped-records-total", "stream-task-metrics", "", tags)); + final Metric dropRate = metrics.get( + new MetricName("dropped-records-rate", "stream-task-metrics", "", tags)); + + assertEquals(3.0, dropTotal.metricValue()); + assertEquals( + 3.0 / 30.0, + ((Number) dropRate.metricValue()).doubleValue(), + 0.005d, + "dropped-records-rate must reflect the 3 records dropped, not the single sensor recording; " + + "counting recordings would give 1/30 == 0.03333 (KAFKA-20877)" + ); + } + } \ No newline at end of file diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/InMemoryWindowStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/InMemoryWindowStoreTest.java index 72af0085de093..82a47420f634f 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/InMemoryWindowStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/InMemoryWindowStoreTest.java @@ -17,6 +17,8 @@ package org.apache.kafka.streams.state.internals; import org.apache.kafka.common.IsolationLevel; +import org.apache.kafka.common.Metric; +import org.apache.kafka.common.MetricName; import org.apache.kafka.common.header.internals.RecordHeaders; import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.serialization.Serdes; @@ -53,6 +55,7 @@ import java.time.Instant; import java.util.LinkedList; import java.util.List; +import java.util.Map; import java.util.Properties; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -167,6 +170,42 @@ public void shouldRestore() { } } + @Test + public void shouldMeasureExpiredRecordsDroppedDuringRestoreAsRecords() { + final StateSerdes serdes = new StateSerdes<>("", Serdes.Integer(), Serdes.String()); + + final List> batch = new LinkedList<>(); + // advances observed stream time far enough that every record after it falls outside retention + batch.add(new KeyValue<>( + toStoreKeyBinary(0, 4 * RETENTION_PERIOD, 0, new RecordHeaders(), serdes).get(), + serdes.rawValue("on-time"))); + for (int key = 1; key <= 3; key++) { + batch.add(new KeyValue<>( + toStoreKeyBinary(key, 0L, 0, new RecordHeaders(), serdes).get(), + serdes.rawValue("expired"))); + } + + context.restore(STORE_NAME, batch); + + final Map metrics = context.metrics().metrics(); + final String threadId = Thread.currentThread().getName(); + final Map tags = mkMap(mkEntry("thread-id", threadId), mkEntry("task-id", "0_0")); + + final Metric dropTotal = metrics.get( + new MetricName("dropped-records-total", "stream-task-metrics", "", tags)); + final Metric dropRate = metrics.get( + new MetricName("dropped-records-rate", "stream-task-metrics", "", tags)); + + assertEquals(3.0, dropTotal.metricValue()); + assertEquals( + 3.0 / 30.0, + ((Number) dropRate.metricValue()).doubleValue(), + 0.005d, + "dropped-records-rate must reflect the 3 records dropped, not the single sensor recording; " + + "counting recordings would give 1/30 == 0.03333 (KAFKA-20877)" + ); + } + @Test public void shouldNotExpireFromOpenIterator() { windowStore.put(1, "one", 0L);