apache / apache/fluss

[flink] currentFetchEventTimeLag is misleading when subscribing multiple partitions/buckets with uneven lag

Open
#3,304 2 comments 1 reaction 0 assignees View on GitHub
Dominant language
Java
Stars
2.1k
Forks
625
Avg merge
3d 14h
Merged PRs (30d)
97

Description

### Search before asking

- [x] I searched in the [issues](https://github.com/apache/fluss/issues) and found nothing similar.

### Fluss version

0.9.0 (latest release)

### Please describe the bug 🐞

When a Flink source subscribes to multiple partition tables (or multiple buckets) whose consuming progress differs significantly, the reported `currentFetchEventTimeLag` metric is much smaller than the actual `currentEmitEventTimeLag`. This makes the metric misleading for monitoring and alerting, because:

- `currentFetchEventTimeLag` stays near 0, suggesting the source has caught up.
- `currentEmitEventTimeLag` reports a large lag (e.g. 2 hours), suggesting records are piling up.
- Yet the Flink job has **no backpressure**, so the gap cannot be explained by downstream slowness.

## Root Cause

In [`FlinkSourceSplitReader#forLogRecords`](file:///Users/loserwang/workstation/github/fluss/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java#L426-L494):

```java
maxConsumerRecordTimestampInFetch =
Math.max(maxConsumerRecordTimestampInFetch, lastRecord.timestamp());
...
flinkSourceReaderMetrics.reportRecordEventTime(
fetchTimestamp - maxConsumerRecordTimestampInFetch);
```

The aggregated timestamp is the **MAX** across all buckets in the fetch, so the reported lag is effectively the **MIN** lag across buckets.

### Example scenario

- Partition A is 2 hours behind → its last record timestamp is 2 hours old.
- Partition B is reading the latest data → its last record timestamp ≈ `now`.

Every `logScanner.poll(POLL_TIMEOUT)` tends to return partition B's fresh data (since it is always locally ready as a `CompletedFetch`). The `Math.max` across buckets then picks partition B's near-`now` timestamp, so:

```
fetchTimestamp - maxTimestamp ≈ 0 → currentFetchEventTimeLag reports ~0
```

Meanwhile, emit lag is computed **per record** and aggregated as **max**, so partition A's 2-hour-old records correctly inflate `currentEmitEventTimeLag`. This is why the two metrics diverge dramatically even though there is no backpressure.

### Solution

1. **Align fetch lag semantics with emit lag (max-lag across buckets):**
Change the aggregation in `forLogRecords` from `Math.max` to `Math.min` on record timestamp (equivalent to `max` on lag), so `currentFetchEventTimeLag` reports the worst-case lag across buckets in the fetch.

2. **Expose per-bucket / per-split fetch lag:**
Add a new gauge `currentFetchEventTimeLag` under the existing per-bucket metric group (next to `currentOffset`):
- Non-partitioned: `fluss.reader.bucket.{bucket_id}.currentFetchEventTimeLag`
- Partitioned: `fluss.reader.partition.{partition_id}.bucket.{bucket_id}.currentFetchEventTimeLag`

This allows users to observe exactly which partition/bucket is lagging, making diagnosis like the scenario above straightforward.

### Are you willing to submit a PR?

- [x] I'm willing to submit a PR!

Contributor guide

No contributing guide indexed for this repository

Research direction

Start in fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java, especially FlinkSourceSplitReader#forLogRecords, and inspect the existing per-bucket currentOffset metric registration. Verify the fetch aggregation reports the worst bucket lag and expose the corresponding per-bucket gauges for partitioned and non-partitioned sources.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
observability, stream-processing
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Quiet
Clarity
Clearly specified
Newbie friendliness
64/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.