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]

Reply via email to