From 1c7d8146df4e0f30827939bedbf7bd9b256bea50 Mon Sep 17 00:00:00 2001 From: "Matthias J. Sax" Date: Fri, 31 Jul 2026 15:14:07 -0700 Subject: [PATCH 1/4] KAFKA-20877: Fix incorrect restore-rate and update-rate metrics Both metrics incorrectly report "batches restored per second" instead of "records restored per seconds". --- .../internals/metrics/TaskMetrics.java | 29 ++++++++++--------- .../processor/internals/StandbyTaskTest.java | 20 +++++++++++-- .../processor/internals/StreamTaskTest.java | 18 ++++++++++-- .../internals/metrics/TaskMetricsTest.java | 13 ++------- 4 files changed, 50 insertions(+), 30 deletions(-) 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/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..4006db03aa3b3 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,9 +75,9 @@ 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.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.doNothing; import static org.mockito.Mockito.doThrow; @@ -487,12 +487,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 + assertTrue( + // regression test for KAFKA-20877: previously we did incorrectly count batches which would result in 0.03333 + // using 0.5 as good intermediate to the expected value of 0.83333 + // -> avoid equalTo(...) on floating point numbers + 0.5d < ((Number) rateMetric.metricValue()).doubleValue(), + "Expected a value larger 0.5 [precisely 0.83333...], but got " + rateMetric.metricValue() + ); task.recordRestoration(time, 50L, false); assertThat(totalMetric.metricValue(), equalTo(75.0)); - assertThat(rateMetric.metricValue(), not(0.0)); + assertTrue( + // regression test for KAFKA-20877: previously we did incorrectly count batches which would result in 0.06666 + // using 2.0 as good intermediate to the expected value of 2.5 + // -> avoid equalTo(...) on floating point numbers + 2.0d < ((Number) rateMetric.metricValue()).doubleValue(), + "Expected a value larger 2.0 [precisely 2.5], but got " + rateMetric.metricValue() + ); } 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..ae63e8008c6a9 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 + assertTrue( + // regression test for KAFKA-20877: previously we did incorrectly count batches which would result in 0.03333 + // using 0.5 as good intermediate to the expected value of 0.83333 + // -> avoid equalTo(...) on floating point numbers + 0.5d < ((Number) rateMetric.metricValue()).doubleValue(), + "Expected a value larger 0.5 [precisely 0.83333...], but got " + rateMetric.metricValue() + ); assertThat(remainMetric.metricValue(), equalTo(75.0)); task.recordRestoration(time, 50L, false); assertThat(totalMetric.metricValue(), equalTo(75.0)); - assertThat(rateMetric.metricValue(), not(0.0)); + assertTrue( + // regression test for KAFKA-20877: previously we did incorrectly count batches which would result in 0.06666 + // using 2.0 as good intermediate to the expected value of 2.5 + // -> avoid equalTo(...) on floating point numbers + 2.0d < ((Number) rateMetric.metricValue()).doubleValue(), + "Expected a value larger 2.0 [precisely 2.5], but got " + rateMetric.metricValue() + ); 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 ) ); From 32f699b06461277ed2bb8ed88ae2979891013606 Mon Sep 17 00:00:00 2001 From: "Matthias J. Sax" Date: Mon, 3 Aug 2026 16:43:01 -0700 Subject: [PATCH 2/4] review comments --- .../processor/internals/StandbyTaskTest.java | 25 +++++----- .../processor/internals/StreamTaskTest.java | 24 ++++----- ...lSchemaRocksDBSegmentedBytesStoreTest.java | 12 ++++- ...bstractRocksDBSegmentedBytesStoreTest.java | 12 ++++- .../AbstractSessionBytesStoreTest.java | 12 ++++- .../AbstractWindowBytesStoreTest.java | 12 ++++- .../internals/InMemorySessionStoreTest.java | 50 +++++++++++++++++++ .../internals/InMemoryWindowStoreTest.java | 42 ++++++++++++++++ 8 files changed, 156 insertions(+), 33 deletions(-) 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 4006db03aa3b3..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 @@ -77,7 +77,6 @@ import static org.hamcrest.Matchers.isA; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; -import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.doNothing; import static org.mockito.Mockito.doThrow; @@ -489,23 +488,23 @@ public void shouldRecordRestoredRecords() { assertThat(totalMetric.metricValue(), equalTo(25.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 - assertTrue( - // regression test for KAFKA-20877: previously we did incorrectly count batches which would result in 0.03333 - // using 0.5 as good intermediate to the expected value of 0.83333 - // -> avoid equalTo(...) on floating point numbers - 0.5d < ((Number) rateMetric.metricValue()).doubleValue(), - "Expected a value larger 0.5 [precisely 0.83333...], but got " + rateMetric.metricValue() + 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)); - assertTrue( - // regression test for KAFKA-20877: previously we did incorrectly count batches which would result in 0.06666 - // using 2.0 as good intermediate to the expected value of 2.5 - // -> avoid equalTo(...) on floating point numbers - 2.0d < ((Number) rateMetric.metricValue()).doubleValue(), - "Expected a value larger 2.0 [precisely 2.5], but got " + rateMetric.metricValue() + 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)" ); } 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 ae63e8008c6a9..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 @@ -895,24 +895,24 @@ public void shouldRecordRestoredRecords() { assertThat(totalMetric.metricValue(), equalTo(25.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 - assertTrue( - // regression test for KAFKA-20877: previously we did incorrectly count batches which would result in 0.03333 - // using 0.5 as good intermediate to the expected value of 0.83333 - // -> avoid equalTo(...) on floating point numbers - 0.5d < ((Number) rateMetric.metricValue()).doubleValue(), - "Expected a value larger 0.5 [precisely 0.83333...], but got " + rateMetric.metricValue() + 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)); - assertTrue( - // regression test for KAFKA-20877: previously we did incorrectly count batches which would result in 0.06666 - // using 2.0 as good intermediate to the expected value of 2.5 - // -> avoid equalTo(...) on floating point numbers - 2.0d < ((Number) rateMetric.metricValue()).doubleValue(), - "Expected a value larger 2.0 [precisely 2.5], but got " + rateMetric.metricValue() + 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/state/internals/AbstractDualSchemaRocksDBSegmentedBytesStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractDualSchemaRocksDBSegmentedBytesStoreTest.java index 3e51e20a35d43..04077f5e09f95 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; @@ -1622,7 +1621,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..bc9a224780669 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; @@ -998,7 +997,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..86a45b5b9aa65 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; @@ -931,7 +930,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..61ff5636accb1 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 @@ -67,7 +67,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; @@ -1044,7 +1043,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..cf9dc64a1263f 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,13 @@ 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.common.utils.Time; +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 +41,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 +324,48 @@ 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. + // + // Align the context's cached system time with the metrics clock, so the rate's sampling window is + // the un-elapsed default of (metrics.num.samples - 1) * metrics.sample.window.ms == 30s. + context.setSystemTimeMs(Time.SYSTEM.milliseconds()); + + 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..b888df91100bb 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,10 +17,13 @@ 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; import org.apache.kafka.common.utils.Bytes; +import org.apache.kafka.common.utils.Time; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.kstream.Windowed; @@ -53,6 +56,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 +171,44 @@ public void shouldRestore() { } } + @Test + public void shouldMeasureExpiredRecordsDroppedDuringRestoreAsRecords() { + context.setSystemTimeMs(Time.SYSTEM.milliseconds()); + + 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); From 4e4ca39d33bd39fbf7cda2b3e62fd37459828bbd Mon Sep 17 00:00:00 2001 From: "Matthias J. Sax" Date: Mon, 3 Aug 2026 17:39:19 -0700 Subject: [PATCH 3/4] align code style --- .../apache/kafka/streams/processor/internals/StreamThread.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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); } From 1c6552f9b4bb5046b3652dd244d2fc5ff596ced7 Mon Sep 17 00:00:00 2001 From: "Matthias J. Sax" Date: Mon, 3 Aug 2026 18:07:46 -0700 Subject: [PATCH 4/4] more cleanup --- .../AbstractDualSchemaRocksDBSegmentedBytesStore.java | 2 +- .../internals/AbstractRocksDBSegmentedBytesStore.java | 2 +- .../streams/state/internals/InMemorySessionStore.java | 8 ++++---- .../streams/state/internals/InMemoryWindowStore.java | 4 ++-- .../streams/state/internals/RocksDBVersionedStore.java | 4 ++-- .../AbstractDualSchemaRocksDBSegmentedBytesStoreTest.java | 2 -- .../internals/AbstractRocksDBSegmentedBytesStoreTest.java | 2 -- .../state/internals/AbstractSessionBytesStoreTest.java | 2 -- .../state/internals/AbstractWindowBytesStoreTest.java | 3 --- .../streams/state/internals/InMemorySessionStoreTest.java | 5 ----- .../streams/state/internals/InMemoryWindowStoreTest.java | 3 --- 11 files changed, 10 insertions(+), 27 deletions(-) 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/state/internals/AbstractDualSchemaRocksDBSegmentedBytesStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractDualSchemaRocksDBSegmentedBytesStoreTest.java index 04077f5e09f95..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 @@ -1585,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 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 bc9a224780669..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 @@ -961,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 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 86a45b5b9aa65..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 @@ -893,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 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 61ff5636accb1..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; @@ -1006,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); 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 cf9dc64a1263f..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 @@ -22,7 +22,6 @@ 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.common.utils.Time; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.kstream.Windowed; @@ -330,10 +329,6 @@ 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. - // - // Align the context's cached system time with the metrics clock, so the rate's sampling window is - // the un-elapsed default of (metrics.num.samples - 1) * metrics.sample.window.ms == 30s. - context.setSystemTimeMs(Time.SYSTEM.milliseconds()); final List> batch = new LinkedList<>(); // advances observed stream time far enough that every record after it falls outside retention 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 b888df91100bb..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 @@ -23,7 +23,6 @@ import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.common.utils.Bytes; -import org.apache.kafka.common.utils.Time; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.kstream.Windowed; @@ -173,8 +172,6 @@ public void shouldRestore() { @Test public void shouldMeasureExpiredRecordsDroppedDuringRestoreAsRecords() { - context.setSystemTimeMs(Time.SYSTEM.milliseconds()); - final StateSerdes serdes = new StateSerdes<>("", Serdes.Integer(), Serdes.String()); final List> batch = new LinkedList<>();