Skip to content

[WIP] Split partition-level bucket-count support into focused PRs - #9370

Open
dwangatt wants to merge 29 commits into
apache:masterfrom
atlassian-forks:dwang/continue-pr-7865
Open

dwangatt wants to merge 29 commits into
apache:masterfrom
atlassian-forks:dwang/continue-pr-7865

Conversation

@dwangatt

@dwangatt dwangatt commented Aug 24, 2026 •

Copy link
Copy Markdown
Contributor

This PR continues #7865.

Support per-partition bucket counts so individual partitions can be rescaled independently without changing the bucket count of the entire table.

This PR is being restructured into a focused series following review feedback. The original change mixed independent test cleanup, core partition-layout semantics, write-path changes, Flink integration, Spark compatibility policy, and documentation.

The independent changes have been moved into separate PRs:

The remaining feature will be rebuilt on top of the core prerequisite PRs:

  1. Make partition-aware overwrite routing a core BatchWriteBuilder contract.
  2. Add the Flink integration and end-to-end coverage.
  3. Explicitly reject unsupported Spark writes.
  4. Document the completed safety boundary.

This PR is kept as WIP while the focused prerequisites are reviewed and merged. It should not be merged in its current form.

Continuation of #7865

The original author is no longer available to maintain the source branch. This continuation preserves all 19 original commits and their authorship, and moves ongoing maintenance to the atlassian-forks/paimon fork.

Original PR:

Purpose

Support per-partition bucket counts so individual partitions can be rescaled independently without changing the bucket count of the entire table.

This allows a table to contain partitions with different bucket layouts while ensuring readers, writers, compaction, and overwrite/rescale operations use the bucket count recorded for each partition.

Commit history

All original commits from #7865 are preserved. New changes addressing outstanding review feedback will be added as separate commits on this maintained branch.

Follow-up

  • Rebase and resolve conflicts against the latest apache/master
  • Review and address unresolved feedback from [core][flink] Support Per-Partition Bucket Counts #7865
  • Run the full unit and integration test suite
  • Verify streaming writer behavior after a partition rescale
  • Verify INSERT OVERWRITE/rescale behavior
  • Verify writer coordinator and empty-bucket restore behavior

Supersedes #7865.

@dwangatt

Copy link
Copy Markdown
Contributor Author

hey @JingsongLi following up the conversion from #7865
I added code so that:

  • when bucket.per-partition-count-enabled is enabled, Paimon sink can use bucket from partitions so partitions can have different buckets
  • when bucket.per-partition-count-enabled is disabled, Paimon sink will use original code.

Can you also re-run the failed action? It failed for some network issue when downloading JAR from maven repository. Re-run should fix this issue, but I do not have permission to trigger it.

cc @mikedias

@JingsongLi

Copy link
Copy Markdown
Contributor

Spark's fast-write path still relies on the table-level coreOptions.bucket() setting rather than routing based on partition-to-bucket mappings; this can lead to errors or data being written to the incorrect bucket when partitions have varying bucket counts. The documentation itself acknowledges a risk of data corruption—an issue that cannot be mitigated simply by following usage guidelines. Another edge case arises in PartitionEntry.merge: when creation times are identical, the logic defaults to the larger bucket count, potentially causing the old layout to be retained during downscaling operations.

@dwangatt
dwangatt force-pushed the dwang/continue-pr-7865 branch 2 times, most recently from 76d952e to 8115f87 Compare September 2, 2026 04:27
@dwangatt
dwangatt force-pushed the dwang/continue-pr-7865 branch from 8115f87 to bbfa37e Compare September 2, 2026 06:19
@dwangatt

dwangatt commented Sep 2, 2026

Copy link
Copy Markdown
Contributor Author

Thanks @JingsongLi I updated PR with:

  • throw an exception for Spark connector when per-partition option is on. It only supports FLINK connector only.
  • check buckets from new and old file, and if two files with same timestamp but buckets are not same, then fail the job. In this case, it should not happen.

}

def write(data: DataFrame): Seq[CommitMessage] = {
if (coreOptions.bucketPerPartitionCountEnabled()) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] This guard only covers the V1 PaimonSparkWriter path. With spark.paimon.write.use-v2-write=true, PaimonSparkTableBase.newWriteBuilder returns PaimonV2WriteBuilder, whose PaimonV2Write.toBatch writes directly through BatchWriteBuilder and never calls this method. Therefore Spark V2 INSERT/OVERWRITE is still accepted even though the docs say Spark rejects this table mode; its PaimonWriteRequirement also still clusters with the table-level bucket count, so a rescaled partition can be distributed with the wrong layout. Please reject this option in the shared Spark write-builder/table entry point (or in both V1 and V2 paths) and cover both write.use-v2-write settings in the test.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@JingsongLi good point. Updated PR so that in v2 API, same exception will be thrown when bucket count is different

case HASH_FIXED:
return new FixedBucketRowKeyExtractor(schema());
return new FixedBucketRowKeyExtractor(
schema(), PartitionBucketMapping.loadFromTable(this));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] This turns every createRowKeyExtractor() into a full partition-manifest scan. In the Flink sink this method is called once on the client for RowDataChannelComputer and again inside every writer subtask when StoreSinkWriteImpl calls table.newWrite; each writer also performs another full PartitionBucketMapping.loadFromScan in FileSystemWriteRestore (and that restore is constructed even when the coordinator later replaces it). With P writer subtasks, startup is therefore roughly 2P full scans/partition-map builds, not the single extra scan described by the option/docs. On the large partitioned tables this feature targets, that can create an object-store request storm and duplicate the full mapping in every task. Please build the mapping once per table/job and pass/share it with both routing and restore, or make the writer-side lookup partition-scoped.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@JingsongLi I added an extra parameter to related functions so that this mapping can be passed to them instead of loading them multiple times.

Comment on lines +682 to +695
if (bucket >= restoredTotalBuckets) {
throw new RuntimeException(
String.format(
"Trying to write bucket %d to %s, but the partition only has %d "
+ "buckets (table default: %d). Recompute the bucket using the "
+ "partition's bucket count, or rescale the partition via "
+ "INSERT OVERWRITE.",
bucket,
partInfo.get(),
restoredTotalBuckets,
expectedTotalBuckets));
}
LOG.info(
"{} uses {} buckets (expected: {}). Accepting per-partition bucket count.",

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] The relaxed bucket check silently accepts mis-routed rows on the Flink path.

bucket >= restoredTotalBuckets is a range check, not a routing check, and it cannot detect the failure this feature makes reachable.

The routing mapping is frozen at job submission: FlinkSinkBuilder builds new RowDataChannelComputer(sinkTable.createRowKeyExtractor()) on the client, and RowKeyExtractor was made Serializable in this PR precisely so the PartitionBucketMapping ships with the job graph. A long-running streaming job therefore keeps routing with the bucket layout that existed at submission time, and the mapping is never refreshed afterwards.

Consider a partition rescaled from 4 to 8 buckets while a streaming job is running. The router still computes h % 4 = b, which is in [0, 4); the true bucket is h % 8, which is either b or b + 4. Since b < 8, the range check passes for every row, and roughly half of them are written into the wrong bucket with only a LOG.info. On a primary-key table the same key then lives in two buckets — duplicates and lost updates, with no error surfaced at any layer. (Downscaling is caught only accidentally, and only for power-of-two counts; e.g. 6 -> 4 also mis-routes silently for h = 8.)

The writer cannot currently do better, because expectedTotalBuckets here is the table-level count: the Flink fixed-bucket path goes TableWriteImpl.writeAndReturn(row, bucket, null) -> write(partition, bucket, data) -> createWriterContainer(partition, bucket, numBuckets, ...), where numBuckets = options.bucket(). The bucket count the router actually used never reaches the writer, so there is nothing meaningful to compare restoredTotalBuckets against.

Suggestion: plumb the routing bucket count (PartitionBucketMapping.resolveNumBuckets(partition)) through to the writer — the write(partition, bucket, totalBuckets, data) overload already exists — and fail when it disagrees with restoredTotalBuckets. That is precise rather than heuristic, and it does not break the feature: a batch job or a freshly started streaming job loads a current mapping, so the two agree and nothing fails. It fails only when the mapping is stale, which is exactly the corruption case that docs/docs/maintenance/rescale-bucket.md currently addresses with "Streaming jobs must be restarted after rescaling a partition" — a guideline that cannot be enforced, as @JingsongLi already noted above.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @zhuxiangyi that's a good point. I updated PR to fail job when bucket does not match.

@dwangatt
dwangatt force-pushed the dwang/continue-pr-7865 branch 2 times, most recently from 85c9a13 to 4927122 Compare September 7, 2026 00:38
@dwangatt
dwangatt force-pushed the dwang/continue-pr-7865 branch from 4927122 to e6c0d6d Compare September 7, 2026 01:48
@dwangatt
dwangatt force-pushed the dwang/continue-pr-7865 branch from 7889585 to 017527d Compare September 7, 2026 03:32
@dwangatt
dwangatt force-pushed the dwang/continue-pr-7865 branch from 017527d to a6255e4 Compare September 7, 2026 05:13
@JingsongLi

JingsongLi commented Sep 20, 2026 •

Copy link
Copy Markdown
Contributor

Thanks for taking over and continuing this work. After another pass, I think we should pause the line-by-line review and restructure this PR first.

The use case is valid, but the current PR boundary is too broad: 49 files and +2,333/-209 lines, spanning global PartitionEntry aggregation semantics, new core write/restore APIs, Flink routing and coordinator behavior, Spark compatibility policy, documentation, and several large test suites. It also contains clearly independent changes, including the Kafka topic-cleanup retry, the BUCKET_APPEND_ORDERED test rewrite, and other drive-by test/helper changes.

This is not only a reviewability concern. The cross-layer ownership has left concrete correctness gaps:

  • Core BatchWriteBuilder overwrite still routes using the existing partition mapping; only the Flink path wraps the table with SchemaBucketFileStoreTable. After changing a partition from 2 to 4 buckets, a core overwrite can hash rows modulo 2 while stamping the resulting files with totalBuckets=4. Later modulo-4 writes may then place the same key in a different bucket.
  • PartitionEntry.merge now derives the layout from data-file creationTime. This is wall-clock metadata, not logical commit order. For example, after a 4→8 rescale followed by snapshot rollback, deleted 8-bucket files can still have the newest timestamp, causing the restored 4-bucket partition to be reported as using 8 buckets.
  • The new RescaleActionITCase does not enable bucket.per-partition-count-enabled, so much of its 244 lines does not actually exercise the feature’s enabled routing path. A focused partitioned streaming restart test is also still missing.

Could we split this into a focused series?

  1. Move unrelated Kafka/test cleanup into separate PRs.
  2. Resolve the core partition-layout and rollback semantics in a prerequisite PR with focused tests.
  3. Make overwrite routing a core BatchWriteBuilder contract instead of a Flink-only workaround.
  4. Layer the Flink integration, Spark rejection, documentation, and end-to-end tests on top, exposing the option only once the complete safety boundary is in place.

Please also clean up the merge/no-op/whitespace commits when restructuring; attribution can still be preserved.

I support the feature goal, but I do not think the current combined PR can be reviewed or merged safely as-is.

@dwangatt dwangatt changed the title [core][flink] Support Per-Partition Bucket Counts [WIP] Split partition-level bucket-count support into focused PRs Sep 21, 2026
@dwangatt

Copy link
Copy Markdown
Contributor Author

Thanks @JingsongLi. I raised three PRs:

@JingsongLi JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for splitting the work. Independently rescaling partitions has clear end-to-end value, and the dependency order in the updated description matches the requested direction.

The prerequisites are now in different states: #10049 is merged, #10048 was closed without merging, and #10052 remains open. Please update the split-PR status in the description accordingly; the independent Kafka cleanup should not gate the bucket-layout work.

Let's keep this PR as the WIP tracking point while the active-file layout and core BatchWriteBuilder overwrite contracts are addressed in focused PRs. The implementation concerns in the previous review remain for those focused PRs to address. The final integration should demonstrate rescale, rollback and a partitioned streaming restart with the option enabled, together with the Spark rejection boundary.

This is a follow-up on the restructuring response at 15916bb2d1; no code changed since the previous review, and I have not rerun the integration suites in this pass.

@dwangatt

Copy link
Copy Markdown
Contributor Author

Raised another PR to route bucket when writing: #10170

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.

4 participants