Skip to content

[flink] Support safe partition-level bucket counts - #10297

Open
dwangatt wants to merge 8 commits into
apache:masterfrom
atlassian-forks:dwang/partition-level-bucket-write
Open

dwangatt wants to merge 8 commits into
apache:masterfrom
atlassian-forks:dwang/partition-level-bucket-write

Conversation

@dwangatt

Copy link
Copy Markdown
Contributor

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

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

@dwangatt
dwangatt marked this pull request as ready for review September 28, 2026 23:46
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants