Xiaobing Fang created FLINK-40575:
-------------------------------------
Summary: Support dynamic table unsubscription and re-subscription
in Fluss pipeline source
Key: FLINK-40575
URL: https://issues.apache.org/jira/browse/FLINK-40575
Project: Flink
Issue Type: Improvement
Reporter: Xiaobing Fang
### Background
`TableDiscoverer` periodically returns the complete current set of subscribed
tables. However, the Fluss pipeline source currently handles only newly
discovered tables.
When a table disappears from the discovery result, `FlussSourceEnumerator`
retains its assigned splits. The source therefore continues reading the table,
and the retained split state may restore it after a checkpoint or savepoint
recovery.
### Expected behavior
The Fluss pipeline source should treat every successful discovery result as the
authoritative current subscription:
- Removing a table from the result stops its records without restarting the job.
- Other subscribed tables continue running normally.
- Restoring from a checkpoint must not resume a table that is still
unsubscribed.
- Re-adding a table starts a new source lifecycle according to
`scan.startup.mode` and emits a new `CreateTableEvent`.
- Removal affects source-side splits and cached metadata only. It must not emit
a `DropTableEvent` or delete the sink table.
- A successful empty result unsubscribes all tables. Discovery failures fail
the job without applying an incomplete result.
### Proposed change
- Detect removals by logical `TablePath`, independently of partition changes.
- Persist pending-removal tombstones in enumerator checkpoint state.
- Block restored splits until a fresh authoritative subscription snapshot is
available.
- Coordinate split cleanup between the enumerator and every source-reader
subtask.
- Unsubscribe active Fluss buckets and remove finished splits through the
normal source-reader cleanup path.
- Clear table-level schema and deserialization caches so re-added tables are
initialized as new tables.
- Prevent pending assignments, failed-reader split returns, and delayed
asynchronous initialization from resurrecting removed tables.
- Document the discovery, removal, recovery, and re-subscription semantics.
The enumerator state serializer will be upgraded to version 2, while retaining
support for restoring version 1 state.
### Acceptance criteria
1. With tables A and B subscribed, removing B stops new B records while A
continues, without restarting the job.
2. Restoring an old checkpoint while B remains unsubscribed produces no new B
records.
3. Re-adding B creates fresh splits using the configured `scan.startup.mode`
rather than continuing its removed split.
4. Removing one partition from a still-subscribed partitioned table is not
treated as logical-table unsubscription.
5. Pending assignments, delayed initialization results, and splits returned
during failover cannot revive an unsubscribed table.
6. The behavior is covered for both Flink 1.20 and Flink 2.2.
### Background
`TableDiscoverer` periodically returns the complete current set of subscribed
tables. However, the Fluss pipeline source currently handles only newly
discovered tables.
When a table disappears from the discovery result, `FlussSourceEnumerator`
retains its assigned splits. The source therefore continues reading the table,
and the retained split state may restore it after a checkpoint or savepoint
recovery.
### Expected behavior
The Fluss pipeline source should treat every successful discovery result as the
authoritative current subscription:
- Removing a table from the result stops its records without restarting the job.
- Other subscribed tables continue running normally.
- Restoring from a checkpoint must not resume a table that is still
unsubscribed.
- Re-adding a table starts a new source lifecycle according to
`scan.startup.mode` and emits a new `CreateTableEvent`.
- Removal affects source-side splits and cached metadata only. It must not emit
a `DropTableEvent` or delete the sink table.
- A successful empty result unsubscribes all tables. Discovery failures fail
the job without applying an incomplete result.
### Proposed change
- Detect removals by logical `TablePath`, independently of partition changes.
- Persist pending-removal tombstones in enumerator checkpoint state.
- Block restored splits until a fresh authoritative subscription snapshot is
available.
- Coordinate split cleanup between the enumerator and every source-reader
subtask.
- Unsubscribe active Fluss buckets and remove finished splits through the
normal source-reader cleanup path.
- Clear table-level schema and deserialization caches so re-added tables are
initialized as new tables.
- Prevent pending assignments, failed-reader split returns, and delayed
asynchronous initialization from resurrecting removed tables.
- Document the discovery, removal, recovery, and re-subscription semantics.
The enumerator state serializer will be upgraded to version 2, while retaining
support for restoring version 1 state.
### Acceptance criteria
1. With tables A and B subscribed, removing B stops new B records while A
continues, without restarting the job.
2. Restoring an old checkpoint while B remains unsubscribed produces no new B
records.
3. Re-adding B creates fresh splits using the configured `scan.startup.mode`
rather than continuing its removed split.
4. Removing one partition from a still-subscribed partitioned table is not
treated as logical-table unsubscription.
5. Pending assignments, delayed initialization results, and splits returned
during failover cannot revive an unsubscribed table.
6. The behavior is covered for both Flink 1.20 and Flink 2.2.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)