Support report backlog information to Flink
- 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.
### Motivation
In [FLIP-309](https://cwiki.apache.org/confluence/display/FLINK/FLIP-309%3A+Support+using+larger+checkpointing+interval+when+source+is+processing+backlog#FLIP309:Supportusinglargercheckpointingintervalwhensourceisprocessingbacklog-ProposedChanges), Flink introduced the concept of **ProcessingBacklog** to indicate whether a record should be processed in a **low-latency** or **high-throughput** manner. **ProcessingBacklog** can be used to dynamically adjust a job’s checkpoint interval at runtime.
The **Fluss Source** can set **ProcessingBacklog** by marking offsets at the starting position to identify **backlog** data and report this signal to Flink, enabling Flink to perform relevant optimizations based on it.
### Proposed Changes
- **Add backlog boundary tracking in `FlinkSourceEnumerator`**
- Introduce `backlogOffsets: Map` to store the backlog boundary offset per bucket.
- Add/extend `resetBacklog(...)` to initialize splits and persist the backlog boundary offsets.
- Update `initPartitionedSplits()` to pass `initializeAtBeginning` / backlog boundary offset information into the created splits.
- Update `handleSourceEvent(...)` to handle a new `BacklogFinishEvent` and update the enumerator’s backlog tracking state.
- **Propagate backlog boundary offset via split metadata**
- Extend split types to carry backlog boundary offset:
- `LogSplit`: add `backlogOffset`
- `HybridSnapshotLogSplit`: add `backlogOffset`
- **Introduce a new source event for backlog completion**
- Add `BacklogFinishEvent extends SourceEvent` to signal that a specific `TableBucket` has finished consuming backlog (i.e., reader has reached/passed the backlog boundary offset).
- **Report backlog completion from `FlinkSourceSplitReader`**
- Add new fields:
- `backlogMarkedOffsets: Map`: backlog boundary offsets per bucket (from assigned splits).
- `onlySnapshotBuckets: Set`: buckets with snapshot-only data (no log backlog to track).
- `backlogEventSentTbls: Set`: dedup to ensure `BacklogFinishEvent` is sent once per bucket.
- `context: SourceReaderContext`: used to send `BacklogFinishEvent` to the coordinator.
- Update reader logic to detect backlog completion and emit events:
- `fetch(...)`: hook backlog completion reporting when a bucket completes backlog.
- `subscribeLog(...)` and `forLogRecords(...)`: check current log offsets against `backlogMarkedOffsets`; once the boundary is reached, send `BacklogFinishEvent` (deduplicated).
### Anything else?
_No response_
### Willingness to contribute
- [x] I'm willing to submit a PR!
Contributor guide
No contributing guide indexed for this repository
Research direction
Start by reading FlinkSourceEnumerator and the LogSplit and HybridSnapshotLogSplit types to trace how split metadata is created and passed to FlinkSourceSplitReader. Then follow fetch, subscribeLog, and forLogRecords, along with SourceEvent handling, to understand where backlog completion is detected. Done means backlog boundaries propagate correctly and each bucket reports completion once, including snapshot-only buckets.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- java
- Domain
- distributed-systems, stream-processing
- Issue type
- Feature
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Stale
- Clarity
- Mostly clear
- Newbie friendliness
- 45/100