Skip to content

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

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

dwangatt wants to merge 17 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
@JingsongLi

Copy link
Copy Markdown
Contributor

Production review at 1944e95879. The end-to-end Flink feature is useful, but two issues need to be addressed before merging.

[P1] Keep compaction and recovery out of the ambiguous-routing guard

AbstractFileStoreWrite.java:504 now calls requirePartitionBucketCount() before looking up an existing writer. compact(partition, bucket, ...) and notifyNewFiles(...) also use this overload, so they throw on every partitioned table with the option enabled, even after a valid explicit-count write has already created the writer.

This is reachable through supported Flink paths: GlobalFullCompactionSinkWrite.submitFullCompaction() calls write.compact() for written buckets, and LookupSinkWrite calls it when restoring active buckets. I reproduced the actual full-compaction writer with a real local file table: a mapped write succeeds, then prepareCommit(false, 1) fails with the new UnsupportedOperationException; the same control with the option disabled succeeds. No rescale is required. Thus a full-compaction changelog job fails when compaction triggers, and a lookup job can fail during restart.

Restrict the ambiguous-count rejection to routing writes. Maintenance/recovery operations must reuse a valid existing writer or resolve the actual partition count when creating one. Moving the guard after the cache lookup alone would still leave dedicated compaction and recovery of a missing writer broken. Please cover full-compaction commit, lookup restoration and notifyNewFiles with the option enabled.

[P2] Preserve the promised no-op for unpartitioned tables in Spark

PaimonSparkTableBase.scala:175–178 and PaimonSparkWriter.scala:130–133 check the option without checking whether the table is partitioned. The PR promises that the option remains a no-op for unpartitioned tables, and the core mapping/writer checks already enforce that distinction. However, a normal unpartitioned fixed-bucket table with this property set cannot accept a Spark INSERT.

I exercised actual SQL with spark.paimon.write.use-v2-write=false and true: both unpartitioned INSERT/readback controls pass with the property disabled, and both fail in the new guard with it enabled. Limit the rejection to partitioned tables that actually use this feature, and align the documentation with the advertised no-op.

Verification: 210 core tests, 14 Flink tests and all 10 existing Spark PaimonSinkTest cases passed with normal Maven checks, including the rescale/savepoint/restart scenarios. The failures above were additional targeted reproductions, outside the PR's existing coverage.

@JingsongLi

Copy link
Copy Markdown
Contributor

Additional [P1] Compute routing from the default-normalized row at StoreSinkWriteImpl.java:157–159 (1944e95879).

The new mapped path calls write.getPartition(row) and write.getBucket(row) on the raw input, then passes that bucket explicitly to TableWriteImpl.writeAndReturn. The latter applies column defaults at line 243 before extracting the stored row, but keeps the already supplied bucket. The written value and the bucket can therefore disagree. FixedBucketSink supplies a non-null mapping for ordinary fixed-bucket tables too, so this regression occurs even when bucket.per-partition-count-enabled is absent/disabled.

I reproduced this through actual Flink SQL:

CREATE TABLE T (a INT, b INT) WITH (
  'bucket'='4', 'bucket-key'='b', 'bucket-function.type'='mod');
CALL sys.alter_column_default_value('default.T', 'b', '5');
INSERT INTO T (a) VALUES (1);
SELECT * FROM T;             -- returns [1, 5]
SELECT * FROM T WHERE b=5;   -- incorrectly returns no rows

The null input routes to bucket 0, while the stored default 5 belongs to bucket 1. The filtered query prunes bucket 0 and silently misses the committed row. As a falsification control, restoring only the old StoreSinkWriteImpl.write(row) body makes the same INSERT and both queries pass; the rest of this PR remained in place. The original PR source was restored after the test.

Resolve the partition, bucket and total count from the same default-normalized row used for storage, preferably inside the existing conversion path rather than duplicating default handling before it. Include a real fixed-bucket default-key INSERT/filter regression test; partition defaults also need to participate in count resolution after normalization.

@JingsongLi

Copy link
Copy Markdown
Contributor

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.

@dwangatt

dwangatt commented Oct 6, 2026

Copy link
Copy Markdown
Contributor Author

Thanks @JingsongLi I updated PR :)

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.

3 participants