[
https://issues.apache.org/jira/browse/FLINK-40575?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated FLINK-40575:
-----------------------------------
Labels: pull-request-available (was: )
> 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
> Priority: Major
> Labels: pull-request-available
>
> h3. 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.
> h3. 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.
> h3. 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.
> h3. 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)