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

   ## Summary
   
   This PR is extracted from the original larger PR #9370 and contains the 
end-to-end implementation required to make partition-level bucket counts safe 
and usable in Flink.
   
   This PR delivers end-to-end support for partition-level bucket counts in 
Flink.
   
   It supersedes the closed core-only prerequisite PR #10170. The 
partition-layout scan prerequisite has already merged in #10052; this PR 
combines the remaining core routing contract with the supported engine path, 
Spark policy, documentation, and integration coverage.
   
   ## Behavior
   
   When `bucket.per-partition-count-enabled = true` on a partitioned 
fixed-bucket table:
   
   - Flink loads the active partition-to-bucket mapping when the job starts.
   - Flink routes each row with that partition’s bucket count and propagates 
the same `totalBuckets` value to the writer.
   - The writer validates its routing layout against restored files. A stale 
streaming job that continues after a partition rescale fails before it can 
silently write incorrectly routed data.
   - `INSERT OVERWRITE` routes with the target layout so a rescaled partition’s 
rewritten files are both hashed and stamped with the new bucket count.
   - Legacy core writes that only provide `(partition, bucket)` are rejected 
when the option is enabled, because they cannot prove which bucket count was 
used for routing.
   - Spark rejects writes through both V1 and V2 write paths; Flink is the 
supported engine for this option.
   
   The option remains a no-op for unpartitioned tables, which retain the 
existing single table-level bucket-count invariant.
   
   ## Documentation
   
   - Adds `bucket.per-partition-count-enabled` to generated core configuration 
documentation.
   - Documents independent partition rescaling, Flink-only write support, Spark 
rejection, and the required streaming workflow:
     savepoint stop → rescale/overwrite → restart.
   
   ## Tests
   
   - Core mapping, extractor, writer restore, and commit validation coverage.
   - Flink partitioned PK rescale → overwrite → subsequent write/read coverage.
   - Flink streaming savepoint restart after a partition rescale, including a 
duplicate-key regression assertion.
   - Spark write rejection with both `spark.paimon.write.use-v2-write=false` 
and `true`.
   
   ## Verification
   
   ```text
   mvn -q -pl paimon-flink/paimon-flink-common -am \
     -Dtest=RescaleBucketITCase -DfailIfNoTests=false test
   
   mvn -q -pl paimon-core \
     -Dtest=PartitionBucketMappingTest,FixedBucketWriteSelectorTest,\
   FixedBucketRowKeyExtractorTest,FileSystemWriteRestoreTest,\
   FileStoreCommitTest,AppendOnlySimpleTableTest test
   
   mvn -q -pl paimon-spark/paimon-spark-ut -am \
     -Dtest=PaimonSinkTest -DfailIfNoTests=false test


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