From cb9efa71dab0406577aa194408a24ef8881d646f Mon Sep 17 00:00:00 2001 From: Hongshun Wang Date: Tue, 8 Sep 2026 11:43:45 +0800 Subject: [PATCH] [flink] Support pendingRecords metric for Fluss source. Track log record lag using scanner high watermarks and current fetch offsets, and aggregate the lag across subscribed buckets. --- .../client/metrics/ScannerMetricGroup.java | 9 +++ .../table/scanner/log/AbstractLogScanner.java | 1 + .../table/scanner/log/BucketScanStatus.java | 16 ++++- .../table/scanner/log/LogScannerStatus.java | 8 +++ .../scanner/log/LogFetchCollectorTest.java | 23 +++++++ .../org/apache/fluss/metrics/MetricNames.java | 1 + .../metrics/FlinkSourceReaderMetrics.java | 15 +++++ .../source/reader/FlinkSourceSplitReader.java | 15 ++++- .../metrics/FlinkSourceReaderMetricsTest.java | 21 ++++++ .../reader/FlinkSourceSplitReaderTest.java | 67 +++++++++++++++++++ .../observability/monitor-metrics.md | 6 ++ 11 files changed, 180 insertions(+), 2 deletions(-) diff --git a/fluss-client/src/main/java/org/apache/fluss/client/metrics/ScannerMetricGroup.java b/fluss-client/src/main/java/org/apache/fluss/client/metrics/ScannerMetricGroup.java index bae44d00ed5..279db462c2c 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/metrics/ScannerMetricGroup.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/metrics/ScannerMetricGroup.java @@ -24,6 +24,7 @@ import org.apache.fluss.metrics.CharacterFilter; import org.apache.fluss.metrics.Counter; import org.apache.fluss.metrics.DescriptiveStatisticsHistogram; +import org.apache.fluss.metrics.Gauge; import org.apache.fluss.metrics.Histogram; import org.apache.fluss.metrics.MeterView; import org.apache.fluss.metrics.MetricNames; @@ -103,6 +104,14 @@ public Counter remoteFetchErrorCount() { return remoteFetchErrorCount; } + /** + * Registers the gauge for the number of log records that have not been fetched. It must be + * called at most once, otherwise the duplicated registration will be ignored with a warning. + */ + public void registerRecordsLagGauge(Gauge recordsLagGauge) { + gauge(MetricNames.SCANNER_RECORDS_LAG, recordsLagGauge); + } + public void recordPollStart(long pollStartMs) { this.pollStartMs = pollStartMs; this.timeMsBetweenPoll = lastPollMs != 0L ? pollStartMs - lastPollMs : 0L; diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/AbstractLogScanner.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/AbstractLogScanner.java index 8e16b35ef75..60ac2acbb11 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/AbstractLogScanner.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/AbstractLogScanner.java @@ -67,6 +67,7 @@ protected AbstractLogScanner( this.logScannerStatus = logScannerStatus; this.logFetcher = logFetcher; this.scannerMetricGroup = scannerMetricGroup; + this.scannerMetricGroup.registerRecordsLagGauge(logScannerStatus::recordsLag); } /** diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/BucketScanStatus.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/BucketScanStatus.java index 738e04a49e3..c99ae783875 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/BucketScanStatus.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/BucketScanStatus.java @@ -23,7 +23,7 @@ @Internal class BucketScanStatus { private long offset; // last consumed position - private long highWatermark; // the high watermark from last fetch + private long highWatermark = -1L; // the high watermark from last fetch, -1 if never fetched // TODO add resetStrategy and nextAllowedRetryTimeMs. public BucketScanStatus() { @@ -49,4 +49,18 @@ public void setOffset(Long offset) { public void setHighWatermark(Long highWatermark) { this.highWatermark = highWatermark; } + + /** + * Returns the number of log records that have not been fetched for this bucket, or 0 if the lag + * is unknown, i.e. the offset is still a sentinel offset (like {@link + * LogScanner#EARLIEST_OFFSET}) not resolved by any fetch yet, or no high watermark has been + * returned by the server yet. The high watermark can also be staler than the offset, in which + * case the lag is 0 as well. + */ + long recordsLag() { + if (offset < 0 || highWatermark < 0) { + return 0L; + } + return Math.max(highWatermark - offset, 0L); + } } diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/LogScannerStatus.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/LogScannerStatus.java index c3ea600cda0..8ae1391a9a2 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/LogScannerStatus.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/LogScannerStatus.java @@ -65,6 +65,14 @@ synchronized void updateOffset(TableBucket tableBucket, long offset) { bucketStatus(tableBucket).setOffset(offset); } + synchronized long recordsLag() { + long recordsLag = 0L; + for (BucketScanStatus bucketScanStatus : bucketStatusMap.bucketStatusMap().values()) { + recordsLag += bucketScanStatus.recordsLag(); + } + return recordsLag; + } + synchronized void assignScanBuckets(Map scanBucketAndOffsets) { for (Map.Entry entry : scanBucketAndOffsets.entrySet()) { TableBucket scanBucket = entry.getKey(); diff --git a/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/log/LogFetchCollectorTest.java b/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/log/LogFetchCollectorTest.java index aafad8c0eeb..e00c07e8970 100644 --- a/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/log/LogFetchCollectorTest.java +++ b/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/log/LogFetchCollectorTest.java @@ -236,14 +236,37 @@ void testShouldContinueConsumeSameCompletedFetchAcrossPolls() throws Exception { ScanRecords firstPoll = collector.collectFetch(logFetchBuffer); assertThat(firstPoll.records(tb).size()).isEqualTo(2); assertThat(logScannerStatus.getBucketOffset(tb)).isEqualTo(2L); + assertThat(logScannerStatus.recordsLag()).isEqualTo(8L); assertThat(completedFetch.isConsumed()).isFalse(); ScanRecords secondPoll = collector.collectFetch(logFetchBuffer); assertThat(secondPoll.records(tb).size()).isEqualTo(2); assertThat(logScannerStatus.getBucketOffset(tb)).isEqualTo(4L); + assertThat(logScannerStatus.recordsLag()).isEqualTo(6L); assertThat(completedFetch.isConsumed()).isFalse(); } + @Test + void testRecordsLagAggregation() { + TableBucket initializedBucket = new TableBucket(DATA1_TABLE_ID, 0); + TableBucket uninitializedBucket = new TableBucket(DATA1_TABLE_ID, 1); + Map scanBuckets = new HashMap<>(); + scanBuckets.put(initializedBucket, 2L); + scanBuckets.put(uninitializedBucket, LogScanner.EARLIEST_OFFSET); + + LogScannerStatus scannerStatus = new LogScannerStatus(); + scannerStatus.assignScanBuckets(scanBuckets); + scannerStatus.updateHighWatermark(initializedBucket, 10L); + assertThat(scannerStatus.recordsLag()).isEqualTo(8L); + + scannerStatus.updateHighWatermark(uninitializedBucket, 7L); + scannerStatus.updateOffset(uninitializedBucket, 3L); + assertThat(scannerStatus.recordsLag()).isEqualTo(12L); + + scannerStatus.unassignScanBuckets(Collections.singletonList(initializedBucket)); + assertThat(scannerStatus.recordsLag()).isEqualTo(4L); + } + @Test void testFilteredEmptyResponseAdvancesOffset() { Configuration conf = new Configuration(); diff --git a/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java b/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java index 644b11fd7e7..766e0e9abc3 100644 --- a/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java +++ b/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java @@ -331,6 +331,7 @@ public class MetricNames { public static final String SCANNER_FETCH_LATENCY_MS = "fetchLatencyMs"; public static final String SCANNER_FETCH_RATE = "fetchRequestsPerSecond"; public static final String SCANNER_BYTES_PER_REQUEST = "bytesPerRequest"; + public static final String SCANNER_RECORDS_LAG = "recordsLag"; public static final String SCANNER_REMOTE_FETCH_BYTES_RATE = "remoteFetchBytesPerSecond"; public static final String SCANNER_REMOTE_FETCH_RATE = "remoteFetchRequestsPerSecond"; public static final String SCANNER_REMOTE_FETCH_ERROR_RATE = "remoteFetchErrorPerSecond"; diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/metrics/FlinkSourceReaderMetrics.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/metrics/FlinkSourceReaderMetrics.java index 400f1969fa1..3d86c3eafb9 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/metrics/FlinkSourceReaderMetrics.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/metrics/FlinkSourceReaderMetrics.java @@ -19,6 +19,7 @@ import org.apache.fluss.flink.source.reader.FlinkSourceReader; +import org.apache.flink.metrics.Gauge; import org.apache.flink.metrics.groups.SourceReaderMetricGroup; import org.apache.flink.runtime.metrics.MetricNames; @@ -33,6 +34,8 @@ public class FlinkSourceReaderMetrics { // For currentFetchEventTimeLag metric private volatile long currentFetchEventTimeLag = UNINITIALIZED; + private volatile boolean pendingRecordsGaugeRegistered; + public FlinkSourceReaderMetrics(SourceReaderMetricGroup sourceReaderMetricGroup) { this.sourceReaderMetricGroup = sourceReaderMetricGroup; } @@ -50,6 +53,18 @@ public void reportRecordEventTime(long lag) { currentFetchEventTimeLag = lag; } + /** + * Registers the gauge reporting the number of log records that have not been fetched yet, which + * backs the standard Flink pendingRecords metric. Only the first registration takes effect, so + * that re-created split readers won't register the metric again. + */ + public void registerPendingRecordsGauge(Gauge pendingRecordsGauge) { + if (!pendingRecordsGaugeRegistered) { + sourceReaderMetricGroup.setPendingRecordsGauge(pendingRecordsGauge); + pendingRecordsGaugeRegistered = true; + } + } + public SourceReaderMetricGroup getSourceReaderMetricGroup() { return sourceReaderMetricGroup; } diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java index 1fff4292175..e1f91692038 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java @@ -40,6 +40,8 @@ import org.apache.fluss.lake.source.LakeSplit; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TablePath; +import org.apache.fluss.metrics.Gauge; +import org.apache.fluss.metrics.MetricNames; import org.apache.fluss.predicate.Predicate; import org.apache.fluss.types.RowType; import org.apache.fluss.utils.CloseableIterator; @@ -127,7 +129,9 @@ public FlinkSourceSplitReader( @Nullable LakeSource lakeSource, FlinkSourceReaderMetrics flinkSourceReaderMetrics) { this.flinkMetricRegistry = - new FlinkMetricRegistry(flinkSourceReaderMetrics.getSourceReaderMetricGroup()); + new FlinkMetricRegistry( + flinkSourceReaderMetrics.getSourceReaderMetricGroup(), + Collections.singleton(MetricNames.SCANNER_RECORDS_LAG)); this.connection = ConnectionFactory.createConnection(flussConf, flinkMetricRegistry); this.table = connection.getTable(tablePath); this.tableId = table.getTableInfo().getTableId(); @@ -146,6 +150,15 @@ public FlinkSourceSplitReader( this.stoppingOffsets = new HashMap<>(); this.emptyLogSplits = new HashSet<>(); this.lakeSource = lakeSource; + + @SuppressWarnings("unchecked") + Gauge recordsLagGauge = + (Gauge) + checkNotNull( + flinkMetricRegistry.getFlussMetric(MetricNames.SCANNER_RECORDS_LAG), + "The gauge %s should have been registered by the log scanner.", + MetricNames.SCANNER_RECORDS_LAG); + flinkSourceReaderMetrics.registerPendingRecordsGauge(recordsLagGauge::getValue); } @Override diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/metrics/FlinkSourceReaderMetricsTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/metrics/FlinkSourceReaderMetricsTest.java index f7bb7f2fcf9..78121dbd107 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/metrics/FlinkSourceReaderMetricsTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/metrics/FlinkSourceReaderMetricsTest.java @@ -24,6 +24,7 @@ import org.junit.jupiter.api.Test; import java.util.Optional; +import java.util.concurrent.atomic.AtomicLong; import static org.assertj.core.api.Assertions.assertThat; @@ -49,4 +50,24 @@ void testCurrentFetchEventTimeLag() { flinkSourceReaderMetrics.reportRecordEventTime(18213L); assertThat((long) currentFetchEventTimeLag.get().getValue()).isEqualTo(18213L); } + + @Test + void testPendingRecords() { + MetricListener metricListener = new MetricListener(); + FlinkSourceReaderMetrics flinkSourceReaderMetrics = + new FlinkSourceReaderMetrics( + InternalSourceReaderMetricGroup.mock(metricListener.getMetricGroup())); + + // the metric is not registered until a log scanner provides its records lag + assertThat(metricListener.getGauge(MetricNames.PENDING_RECORDS)).isEmpty(); + + AtomicLong recordsLag = new AtomicLong(10L); + flinkSourceReaderMetrics.registerPendingRecordsGauge(recordsLag::get); + Optional> pendingRecords = metricListener.getGauge(MetricNames.PENDING_RECORDS); + assertThat(pendingRecords).isPresent(); + assertThat((long) pendingRecords.get().getValue()).isEqualTo(10L); + + recordsLag.set(3L); + assertThat((long) pendingRecords.get().getValue()).isEqualTo(3L); + } } diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReaderTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReaderTest.java index 5c7cb313f28..394f2c27a96 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReaderTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReaderTest.java @@ -23,6 +23,8 @@ import org.apache.fluss.client.table.writer.AppendWriter; import org.apache.fluss.client.table.writer.UpsertWriter; import org.apache.fluss.client.write.HashBucketAssigner; +import org.apache.fluss.config.ConfigOptions; +import org.apache.fluss.config.Configuration; import org.apache.fluss.flink.lake.split.LakeSnapshotAndFlussLogSplit; import org.apache.fluss.flink.source.metrics.FlinkSourceReaderMetrics; import org.apache.fluss.flink.source.split.HybridSnapshotLogSplit; @@ -43,7 +45,9 @@ import org.apache.flink.connector.base.source.reader.RecordsWithSplitIds; import org.apache.flink.connector.base.source.reader.splitreader.SplitsAddition; import org.apache.flink.connector.base.source.reader.splitreader.SplitsChange; +import org.apache.flink.metrics.Gauge; import org.apache.flink.metrics.testutils.MetricListener; +import org.apache.flink.runtime.metrics.MetricNames; import org.apache.flink.runtime.metrics.groups.InternalSourceReaderMetricGroup; import org.apache.flink.table.api.ValidationException; import org.junit.jupiter.api.Test; @@ -56,6 +60,7 @@ import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.OptionalLong; import java.util.Set; @@ -245,6 +250,68 @@ void testHandleLogSplitChangesAndFetch() throws Exception { } } + @Test + void testPendingRecordsMetric() throws Exception { + Schema schema = + Schema.newBuilder() + .column("id", DataTypes.INT()) + .column("name", DataTypes.STRING()) + .build(); + TablePath tablePath = TablePath.of(DEFAULT_DB, "test-pending-records-metric"); + long tableId = + createTable( + tablePath, + TableDescriptor.builder().schema(schema).distributedBy(1).build()); + appendRows(tablePath, 5); + + Configuration sourceConf = new Configuration(clientConf); + sourceConf.setInt(ConfigOptions.CLIENT_SCANNER_LOG_MAX_POLL_RECORDS, 1); + MetricListener metricListener = new MetricListener(); + FlinkSourceReaderMetrics sourceReaderMetrics = + new FlinkSourceReaderMetrics( + InternalSourceReaderMetricGroup.mock(metricListener.getMetricGroup())); + + try (FlinkSourceSplitReader splitReader = + new FlinkSourceSplitReader( + sourceConf, + tablePath, + schema.getRowType(), + null, + null, + null, + sourceReaderMetrics)) { + // the metric is registered when the split reader creates the log scanner, and reports + // 0 before anything is fetched + Optional> pendingRecords = + metricListener.getGauge(MetricNames.PENDING_RECORDS); + assertThat(pendingRecords).isPresent(); + assertThat((long) pendingRecords.get().getValue()).isEqualTo(0L); + + TableBucket tableBucket = new TableBucket(tableId, 0); + LogSplit logSplit = new LogSplit(tableBucket, null, 0L); + splitReader.handleSplitsChanges( + new SplitsAddition<>(Collections.singletonList(logSplit))); + + // fetch the rows one by one, the lag should decrease accordingly. Note that a fetch + // may return no records when the poll times out before any record arrives. + int fetchedRows = 0; + while (fetchedRows < 5) { + RecordsWithSplitIds records = splitReader.fetch(); + int rowsInFetch = 0; + if (records.nextSplit() != null) { + while (records.nextRecordFromSplit() != null) { + rowsInFetch++; + } + } + records.recycle(); + if (rowsInFetch > 0) { + fetchedRows += rowsInFetch; + assertThat((long) pendingRecords.get().getValue()).isEqualTo(5 - fetchedRows); + } + } + } + } + @Test void testHandleMixSnapshotLogSplitChangesAndFetch() throws Exception { TablePath tablePath = TablePath.of(DEFAULT_DB, "test-mix-snapshot-log-table"); diff --git a/website/docs/maintenance/observability/monitor-metrics.md b/website/docs/maintenance/observability/monitor-metrics.md index 4d69d66484c..9d8a0031cf8 100644 --- a/website/docs/maintenance/observability/monitor-metrics.md +++ b/website/docs/maintenance/observability/monitor-metrics.md @@ -1258,6 +1258,12 @@ How to Use Flink Metrics, you can see [Flink Metrics](https://nightlies.apache.o Time difference between reading the data file and file creation. Gauge + + pendingRecords + Flink Source Operator + The number of log records that are available after the current source fetch offset. Only the streaming log part is counted, snapshot and lake records are excluded. + Gauge +