bube opened a new pull request, #10109:
URL: https://github.com/apache/paimon/pull/10109

   ### Purpose
   
   Fixes #10107.
   
   `sink.operator-uid.suffix` and `source.operator-uid.suffix` reach only 
Writer, Global Committer, dynamic-bucket-assigner and Source. Every other 
operator a streaming Paimon read or write adds takes its id from the stream 
graph shape, so any topology change orphans its checkpoint entry, and a restore 
from the HA checkpoint store fails with `There is no operator for the state 
<id>`. Discussed on dev@: [\[DISCUSS\] The operator-uid.suffix options do not 
cover every 
operator](https://lists.apache.org/thread/szjmms72lsrktqps37g5j62mcjtklc9o).
   
   This adds `sink.operator-uid.cover-all-operators` and 
`source.operator-uid.cover-all-operators`, boolean, default false. With the 
option and the matching suffix set, every streaming operator gets 
`<prefix>_<table>_<suffix>`. Existing uids and the default graph are unchanged. 
Batch-only operators are out of scope: they never restore from a checkpoint.
   
   Things a reviewer should know:
   
   - **Enabling the option orphans the newly named operators' old entries 
once.** All are empty except `Collect Statistics` (a 
`CoordinatedOperatorFactory`, present under `partition.sink-strategy = 
PARTITION_DYNAMIC`) and the `Monitor` source on the dedicated-split-generation 
and exactly-once-consumer routes. On the explicit-path route those two need 
`execution.state-recovery.ignore-unclaimed-state` once; a recovery from the HA 
checkpoint store cannot skip entries at all, so the migration goes through a 
savepoint or an explicit checkpoint path. The guide in 
`docs/docs/flink/savepoint.md` ("Operator UIDs") spells this out and shows that 
source, writer and committer state stays mapped.
   - **The final `end` sink is covered.** The HA route rejects any unclaimed 
entry regardless of content, `end` acknowledges checkpoints and so has one, and 
SQL users have no way to uid it themselves. It only changes with the option on, 
and a DataStream user's own `.uid()` on the returned sink still wins. Open for 
discussion, as is the option name.
   - **Same-table-name collision.** `db1.orders` and `db2.orders` written with 
one suffix collide on `Writer_orders_s`. This is pre-existing for the four old 
uids and now applies to about twenty; the guide tells users to use a distinct 
suffix per table.
   - **Uid prefixes equal the operator display names** (spaces included) and 
live in `*_NAME` constants, following `WRITER_NAME`, `GLOBAL_COMMITTER_NAME` 
and `SOURCE_NAME`. They are pinned literally in tests.
   
   ### Tests
   
   - `OperatorUidSuffixTest` (17) and `SourceOperatorUidSuffixTest` (12): 12 
sink and 7 source topology shapes each built under 
`pipeline.auto-generate-uids=false`, so Flink itself rejects any uid-less 
operator; exact node-name guards per shape; literal pins for the four 
pre-existing uids and for every new one; option off and option-without-suffix 
leave today's uids only; operator ids survive an upstream topology change.
   - `OperatorUidSuffixRestoreITCase` (2): restore from a retained checkpoint 
store (HA route, `EmbeddedHaServicesWithLeadershipControl` + 
`PerJobCheckpointRecoveryFactory`) after adding an upstream operator succeeds 
with the option on and fails with `There is no operator for the state` with it 
off.
   - `OperatorUidMigrationITCase` (2): checkpoint taken without the option, 
restored with it through `execution.state-recovery.path`. Unaware-bucket 
succeeds. PARTITION_DYNAMIC fails with `Cannot map checkpoint/savepoint state 
for operator`, then succeeds with 
`execution.state-recovery.ignore-unclaimed-state`; the restored job commits 
under the same `commitUser` (committer state mapped) and the dropped entries, 
read from the checkpoint metadata, are exactly the previously uid-less 
operators.
   - Full `paimon-flink-common` unit phase on Flink 1.20.1 (1094 tests), the 
four classes above on Flink 2.2.0 (`-Pflink2`), `paimon-flink-cdc` compile, 
`ConfigOptionsDocsCompletenessITCase`, and `-DskipTests verify` (checkstyle, 
spotless, rat, enforcer). Not built: the `paimon-flink-1.16..2.1` version 
modules and `paimon-flink-action`.
   
   🤖 Generated with [Claude Code](https://claude.com/claude-code)
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to