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)

Reply via email to