Conversation
…estart note for per-partition bucket rescaling
|
hey @JingsongLi following up the conversion from #7865
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 |
|
Spark's fast-write path still relies on the table-level |
76d952e to
8115f87
Compare
8115f87 to
bbfa37e
Compare
|
Thanks @JingsongLi I updated PR with:
|
| } | ||
|
|
||
| def write(data: DataFrame): Seq[CommitMessage] = { | ||
| if (coreOptions.bucketPerPartitionCountEnabled()) { |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
@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)); |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
@JingsongLi I added an extra parameter to related functions so that this mapping can be passed to them instead of loading them multiple times.
| 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.", |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
Thanks @zhuxiangyi that's a good point. I updated PR to fail job when bucket does not match.
85c9a13 to
4927122
Compare
4927122 to
e6c0d6d
Compare
7889585 to
017527d
Compare
017527d to
a6255e4
Compare
|
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:
Could we split this into a focused series?
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. |
|
Thanks @JingsongLi. I raised three PRs:
|
JingsongLi
left a comment
There was a problem hiding this comment.
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.
|
Raised another PR to route bucket when writing: #10170 |
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:
BatchWriteBuildercontract.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/paimonfork.Original PR:
mdias/master/buckets-per-partitionPurpose
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
apache/masterSupersedes #7865.