[ 
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)

Reply via email to