JingsongLi commented on PR #10297: URL: https://github.com/apache/paimon/pull/10297#issuecomment-5966064435
Reviewed `ce7c4944c76b5f9c977fbf0bbcad8dff22313fc5`. The previous blanket maintenance guard, default-value routing, and unpartitioned Spark rejection cases are fixed. Two production cases remain: **[P1] Retire obsolete bucket IDs when restoring maintenance state** (`AbstractFileStoreWrite.java:724–734`, callers `LookupSinkWrite.java:93–100` and `GlobalFullCompactionSinkWrite.java:128–134,267`). I captured real checkpoint state through both sink implementations’ `snapshotState()` calls, committed the checkpoint, and closed the sinks. Initially partitions P1/P2 and the table default all use 4 buckets, with a real key in bucket 3. I then physically overwrote only P1 to 2 buckets using the same dynamic schema-copy target as `RescaleAction.withBucketNum(2)`; P2/default remain 4 and all rows are preserved. Restoring the captured state fails during lookup sink construction, or during the full-compaction sink’s first forced `prepareCommit`, because the retired P1/bucket 3 now trips the new range guard. This blocks the documented stop → rescale → restart workflow when shrinking a partition. Filtering only that retired P1/bucket-3 state value, with the rest of the state and head code unchanged, makes both restores succeed: update/readback, actual COMPACT/changelog publication and totalBuckets=2/4 stamping all pass. Keeping the original state but changing only the writer default to 2 also succeeds, showing that rejection depends on an unrelated default. Please reconcile maintenance state with the current layout and retire obsolete bucket IDs while retaining valid state and strict routing-write checks. Add shrink/restart coverage for lookup and full-compaction producers. The additional reproduction uses actual sink checkpoint-state contents and local commits, not a serialized binary savepoint; the existing MiniCluster savepoint test already passes but does not cover this case. **[P1] Validate an empty bucket against authoritative current partition state** (`FileSystemWriteRestore.java:86–94,123–128`, `WriteRestore.java:63–64`, `StoreSinkWriteImpl.java:125–126`). The new restore path reuses the routing mapping captured when the job was built. When a target bucket is empty, it therefore “validates” the old count against that same old mapping instead of the snapshot being restored. I initialized a real writer/extractor/restore with schema default 4 and P1’s actual count 2, capturing P1→2 before overwrite. After physically rescaling P1 to 4, keys 0 and 3 are in buckets 0 and 3. A stale mapped update of key 3 to value 999 opens the previously unopened bucket 1, passes restoration using count 2, and commits through the ordinary commit path. Manifests then contain both counts 2/4; **both the unfiltered read and key=3 read return `[1,3,300]` and `[1,3,999]`, violating primary-key uniqueness**. The option-disabled control hashes into bucket 3/count 4 and keeps one updated key. An isolated control that uses fresh restore-layout information and retains the count of existing default-count partitions rejects the stale write before publishing. Both details matter: `loadFromScan` currently discards partitions equal to the default, so null also conflates a known default-count partition with an unseen partition. Please resolve the authoritative count for the snapshot used by restoration independently of the cached routing mapping, including empty buckets/default-count partitions. The existing `commit(..., true)` tests do not establish safety of the normal `commit(..., false)` path; that commit-check default predates this PR. This case exercises the advertised protection against a stale/queued job after rescale, rather than changing the required stop/rescale/restart workflow. Validation: **241 normal JDK 8 tests passed** (211 core, 19 Flink 1.20.1, 11 Spark 3), with Checkstyle/Spotless/enforcer enabled. Three additional four-writer Flink SQL scenarios verified mixed default/explicit values, partition overwrite and branch layout/readback. Valid dedicated compaction and empty-bucket notify/compaction also pass. The two findings above come from extra actual-state/file reproductions with controls. The remote Spark 4 job failed in the option-disabled non-bucketed concurrent-MERGE snapshot-expiration test; its failed jobs have been rerun and are still pending, so I am not claiming a fully green remote run. -- 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]
