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]