Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
5cf7b02
Support per-partition buckets
mikedias Mar 2, 2026
f9b0a1f
Optimize PartitionBucketMapping.loadFromTable
mikedias Mar 9, 2026
81f5d44
Fix rescaling via INSERT OVERWRITE
mikedias Apr 15, 2026
012e72a
Fix empty-bucket from WriteRestore scenario
mikedias May 14, 2026
0d363ce
Fixing corner case for non-partitioned tables
mikedias May 15, 2026
9c784bf
Fix merge commutativity, fail fast on scan error, and add streaming r…
mikedias May 25, 2026
2360910
Improve how we test the BUCKET_APPEND_ORDERED behaviour
mikedias May 26, 2026
45f657f
Merge branch 'master' into mdias/master/buckets-per-partition
mikedias Jun 13, 2026
cc5a260
Fix conflicts with TableWriteCoordinator
mikedias Jun 13, 2026
ede2957
Merge branch 'master' into mdias/master/buckets-per-partition
mikedias Jun 20, 2026
c5172e8
Reject bucket writes outside partition layout
mikedias Jun 24, 2026
d06e210
Update docs to reflect the limitations around Spark
mikedias Jun 25, 2026
c37b785
docs tweaking
mikedias Jul 7, 2026
ed95279
merge with master
mikedias Jul 7, 2026
49dbe75
trailing whitespaces
mikedias Jul 7, 2026
e0d5bd1
Add bucket.per-partition-count-enabled config
mikedias Jul 16, 2026
ffcd6ea
merge with master
mikedias Jul 16, 2026
1bd7597
Validate the option before allowing divergent bucket writing
mikedias Jul 16, 2026
5500330
fix ReadWriteTableITCase
mikedias Jul 17, 2026
a0d1240
Merge apache/master and guard partition bucket mismatches
dwangatt Aug 24, 2026
589d275
Preserve FixedBucketWriteSelector constructor compatibility
dwangatt Aug 24, 2026
4062a90
Preserve postpone bucket restore semantics
dwangatt Aug 24, 2026
219235c
Trigger CI
dwangatt Aug 25, 2026
78f06c7
Clarify partition bucket mapping serialization
dwangatt Aug 25, 2026
45d5844
Retry Kafka topic cleanup timeouts in tests
dwangatt Aug 25, 2026
e6c0d6d
fix: reject Spark writes with per-partition buckets
dwangatt Sep 2, 2026
e0808e5
Merge branch 'master' into dwang/continue-pr-7865
dwangatt Sep 7, 2026
a6255e4
fix: reject Spark writes with per-partition buckets
dwangatt Sep 2, 2026
15916bb
Merge branch 'master' into dwang/continue-pr-7865
dwangatt Sep 14, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
58 changes: 44 additions & 14 deletions docs/docs/maintenance/rescale-bucket.md
Original file line number Diff line number Diff line change
Expand Up @@ -45,15 +45,12 @@ Please note that
- `ALTER TABLE` only modifies the table's metadata and will **NOT** reorganize or reformat existing data.
Reorganize existing data must be achieved by `INSERT OVERWRITE`.
- Rescale bucket number does not influence the read and running write jobs.
- Once the bucket number is changed, any newly scheduled `INSERT INTO` jobs which write to without-reorganized
existing table/partition will throw a `TableException` with message like
```text
Try to write table/partition ... with a new bucket num ...,
but the previous bucket num is ... Please switch to batch mode,
and perform INSERT OVERWRITE to rescale current data layout first.
```
- For partitioned table, it is possible to have different bucket number for different partitions. *E.g.*
- For **partitioned tables**, it is possible to have different bucket numbers for different partitions.
This requires setting `'bucket.per-partition-count-enabled' = 'true'` on the table; otherwise every
partition uses the single table-level bucket count. *E.g.*
```sql
ALTER TABLE my_table SET ('bucket.per-partition-count-enabled' = 'true');

ALTER TABLE my_table SET ('bucket' = '4');
INSERT OVERWRITE my_table PARTITION (dt = '2022-01-01')
SELECT * FROM ...;
Expand All @@ -62,7 +59,40 @@ Please note that
INSERT OVERWRITE my_table PARTITION (dt = '2022-01-02')
SELECT * FROM ...;
```
After these operations, partition `dt=2022-01-01` uses 4 buckets, `dt=2022-01-02` uses 8 buckets, and any
new partitions will use the latest table-level default (8 buckets in this case).
Each partition retains its own bucket count from its data files,
and the new bucket count only applies to newly created partitions or partitions that
have been reorganized with `INSERT OVERWRITE`.

:::info
Per-partition bucket counts are disabled by default. Set `'bucket.per-partition-count-enabled' = 'true'`
to let partitions keep their own bucket count and be rescaled independently. When it is disabled, all
partitions share the table-level `bucket` value.

Note that enabling this option adds an extra manifest scan on write (to resolve each partition's bucket
count), so only enable it when you actually need different bucket counts across partitions.
:::

:::warning
Per-partition bucket counts are currently supported by the **Flink** engine only. Spark rejects writes
to a table with `'bucket.per-partition-count-enabled' = 'true'`; use Flink to write such a table.
:::
- **Unpartitioned tables** require a full rescale before writing. If you change the bucket number and attempt
to write without reorganizing the data first, a `RuntimeException` will be thrown:
```text
Try to write table/partition ... with a new bucket num ...,
but the previous bucket num is ... Please switch to batch mode,
and perform INSERT OVERWRITE to rescale current data layout first.
```
- During overwrite period, make sure there are no other jobs writing the same table/partition.
- **Streaming jobs must be restarted after rescaling a partition.** The per-partition bucket mapping
is loaded once when the streaming job starts (from the manifest files at that point in time). If a
partition is rescaled while the streaming job is running, the job will continue routing rows using
the old bucket count for that partition, which can cause rows to land in wrong buckets and lead to
data correctness issues. The recommended workflow is: suspend the streaming job with a savepoint →
perform the rescale overwrite → restart from the savepoint.


## Use Case

Expand Down Expand Up @@ -106,10 +136,10 @@ SELECT trade_order_id,
FROM raw_orders
WHERE order_status = 'verified';
```
The pipeline has been running well for the past few weeks. However, the data volume has grown fast recently,
and the job's latency keeps increasing. To improve the data freshness, users can
- Suspend the streaming job with a savepoint ( see
[Suspended State](https://nightlies.apache.org/flink/flink-docs-stable/docs/internals/job_scheduling/) and
The pipeline has been running well for the past few weeks. However, the data volume has grown fast recently,
and the job's latency keeps increasing. To improve the data freshness, users can
- Suspend the streaming job with a savepoint ( see
[Suspended State](https://nightlies.apache.org/flink/flink-docs-stable/docs/internals/job_scheduling/) and
[Stopping a Job Gracefully Creating a Final Savepoint](https://nightlies.apache.org/flink/flink-docs-stable/docs/deployment/cli/#terminating-a-job) )
```bash
$ ./bin/flink stop \
Expand Down Expand Up @@ -142,8 +172,8 @@ and the job's latency keeps increasing. To improve the data freshness, users can
FROM verified_orders
WHERE dt IN ('2022-06-20', '2022-06-21', '2022-06-22');
```
- After overwrite job has finished, switch back to streaming mode. And now, the parallelism can be increased alongside with bucket number to restore the streaming job from the savepoint
( see [Start a SQL Job from a savepoint](https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/table/sqlclient/#start-a-sql-job-from-a-savepoint) )
- After overwrite job has finished, switch back to streaming mode. And now, the parallelism can be increased alongside with bucket number to restore the streaming job from the savepoint
( see [Start a SQL Job from a savepoint](https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/table/sqlclient/#start-a-sql-job-from-a-savepoint) )
```sql
SET 'execution.runtime-mode' = 'streaming';
SET 'execution.savepoint.path' = <savepointPath>;
Expand Down
5 changes: 5 additions & 0 deletions docs/docs/primary-key-table/data-distribution.md
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,11 @@ split one bucket into independent key ranges, so bucket count is not a hard read
To change an existing layout, use the offline [Rescale Bucket](../maintenance/rescale-bucket)
workflow.

For partitioned tables, each partition can have its own bucket count when
`'bucket.per-partition-count-enabled' = 'true'` is set. In that case, after a rescale operation existing
partitions retain their original bucket count while newly created partitions use the updated table-level
default. When the option is disabled (the default), all partitions share the single table-level bucket count.

## Dynamic Bucket

Dynamic buckets are the default for primary-key tables (`bucket = -1`). Paimon maintains an index
Expand Down
6 changes: 6 additions & 0 deletions docs/generated/core_configuration.html
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,12 @@
<td>String</td>
<td>Specify the paimon distribution policy. Data is assigned to each bucket according to the hash value of bucket-key.<br />If you specify multiple fields, delimiter is ','.<br />If not specified, the primary key will be used; if there is no primary key, the full row will be used.</td>
</tr>
<tr>
<td><h5>bucket.per-partition-count-enabled</h5></td>
<td style="word-wrap: break-word;">false</td>
<td>Boolean</td>
<td>Whether to allow individual partitions of a fixed-bucket table to keep their own bucket count, so that a single partition can be rescaled independently. Enabling this scans the manifest to resolve each partition's bucket count on write.</td>
</tr>
<tr>
<td><h5>cache-page-size</h5></td>
<td style="word-wrap: break-word;">64 kb</td>
Expand Down
14 changes: 14 additions & 0 deletions paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
Original file line number Diff line number Diff line change
Expand Up @@ -156,6 +156,16 @@ public class CoreOptions implements Serializable {
.withDescription(
"Whether to ignore the order of the buckets when reading data from an append-only table.");

public static final ConfigOption<Boolean> BUCKET_PER_PARTITION_COUNT_ENABLED =
key("bucket.per-partition-count-enabled")
.booleanType()
.defaultValue(false)
.withDescription(
"Whether to allow individual partitions of a fixed-bucket table to keep "
+ "their own bucket count, so that a single partition can be rescaled "
+ "independently. Enabling this scans the manifest to resolve each "
+ "partition's bucket count on write.");

@Immutable
public static final ConfigOption<BucketFunctionType> BUCKET_FUNCTION_TYPE =
key("bucket-function.type")
Expand Down Expand Up @@ -3155,6 +3165,10 @@ public int bucket() {
return options.get(BUCKET);
}

public boolean bucketPerPartitionCountEnabled() {
return options.get(BUCKET_PER_PARTITION_COUNT_ENABLED);
}

public BucketFunctionType bucketFunctionType() {
return options.get(BUCKET_FUNCTION_TYPE);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -83,13 +83,47 @@ public int totalBuckets() {
}

public PartitionEntry merge(PartitionEntry entry) {
PartitionEntry newer = entry.lastFileCreationTime >= lastFileCreationTime ? entry : this;
PartitionEntry older = newer == entry ? this : entry;

// Use the totalBuckets from the most recently created file. This correctly handles
// the case where a partition has been overwritten with a different bucket count: the
// newer files carry the new totalBuckets. A creation timestamp does not provide an
// ordering when the files were created in the same millisecond. In that case, choosing
// either bucket count (for example, the larger one) can silently discard a rescale.
if (newer.lastFileCreationTime == older.lastFileCreationTime
&& newer.totalBuckets != older.totalBuckets) {
if (newer.totalBuckets > 0 && older.totalBuckets > 0) {
throw new IllegalStateException(
String.format(
"Cannot determine the bucket layout for partition %s: files created at %s "
+ "have conflicting bucket counts (%s and %s).",
partition,
newer.lastFileCreationTime,
newer.totalBuckets,
older.totalBuckets));
}

// A non-positive bucket count is a mode marker, not a fixed-bucket layout. For
// example, dynamic bucket tables can contain -1 entries while their real buckets
// are being created. Preserve the previous deterministic tie-break for these modes.
return new PartitionEntry(
partition,
recordCount + entry.recordCount,
fileSizeInBytes + entry.fileSizeInBytes,
fileCount + entry.fileCount,
newer.lastFileCreationTime,
Math.max(newer.totalBuckets, older.totalBuckets));
}
int newTotalBuckets = newer.totalBuckets;

return new PartitionEntry(
partition,
recordCount + entry.recordCount,
fileSizeInBytes + entry.fileSizeInBytes,
fileCount + entry.fileCount,
Math.max(lastFileCreationTime, entry.lastFileCreationTime),
entry.totalBuckets);
newer.lastFileCreationTime,
newTotalBuckets);
}

public Partition toPartition(InternalRowPartitionComputer computer) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@
import org.apache.paimon.partition.PartitionTimeExtractor;
import org.apache.paimon.table.sink.CommitMessage;
import org.apache.paimon.table.sink.CommitMessageImpl;
import org.apache.paimon.table.sink.PartitionBucketMapping;
import org.apache.paimon.types.RowType;
import org.apache.paimon.utils.CommitIncrement;
import org.apache.paimon.utils.ExecutorThreadFactory;
Expand Down Expand Up @@ -152,6 +153,15 @@ public FileStoreWrite<T> withWriteRestore(WriteRestore writeRestore) {
return this;
}

@Override
public FileStoreWrite<T> withPartitionBucketMapping(
PartitionBucketMapping partitionBucketMapping) {
if (restore instanceof FileSystemWriteRestore) {
((FileSystemWriteRestore) restore).withPartitionBucketMapping(partitionBucketMapping);
}
return this;
}

@Override
public FileStoreWrite<T> withIOManager(IOManager ioManager) {
this.ioManager = ioManager;
Expand Down Expand Up @@ -503,7 +513,8 @@ private WriterContainer<T> getWriterWrapper(BinaryRow partition, int bucket, int
partition, totalBuckets, buckets.values().iterator().next().totalBuckets);
}
return buckets.computeIfAbsent(
bucket, k -> createWriterContainer(partition.copy(), bucket, totalBuckets));
bucket,
k -> createWriterContainer(partition.copy(), bucket, totalBuckets, true, true));
}

private Map<Integer, WriterContainer<T>> getWriterContainers(BinaryRow partition) {
Expand All @@ -520,16 +531,20 @@ public RecordWriter<T> createWriter(BinaryRow partition, int bucket) {
}

public WriterContainer<T> createWriterContainer(BinaryRow partition, int bucket) {
return createWriterContainer(partition, bucket, numBuckets, !ignoreNumBucketCheck);
return createWriterContainer(partition, bucket, numBuckets, !ignoreNumBucketCheck, false);
}

private WriterContainer<T> createWriterContainer(
BinaryRow partition, int bucket, int totalBuckets) {
return createWriterContainer(partition, bucket, totalBuckets, true);
return createWriterContainer(partition, bucket, totalBuckets, true, true);
}

private WriterContainer<T> createWriterContainer(
BinaryRow partition, int bucket, int expectedTotalBuckets, boolean validateNumBuckets) {
BinaryRow partition,
int bucket,
int expectedTotalBuckets,
boolean validateNumBuckets,
boolean strictBucketCount) {
if (LOG.isDebugEnabled()) {
LOG.debug("Creating writer for partition {}, bucket {}", partition, bucket);
}
Expand All @@ -554,7 +569,11 @@ private WriterContainer<T> createWriterContainer(
if (!actualIgnorePreviousFiles) {
restored =
scanExistingFileMetas(
partition, bucket, expectedTotalBuckets, validateNumBuckets);
partition,
bucket,
expectedTotalBuckets,
validateNumBuckets,
strictBucketCount);
}

DynamicBucketIndexMaintainer indexMaintainer =
Expand Down Expand Up @@ -647,7 +666,11 @@ public FileStoreWrite<T> withMetricRegistry(MetricRegistry metricRegistry) {
}

private RestoreFiles scanExistingFileMetas(
BinaryRow partition, int bucket, int expectedTotalBuckets, boolean validateNumBuckets) {
BinaryRow partition,
int bucket,
int expectedTotalBuckets,
boolean validateNumBuckets,
boolean strictBucketCount) {
Supplier<String> partInfo =
() ->
partitionType.getFieldCount() > 0
Expand All @@ -674,8 +697,33 @@ private RestoreFiles scanExistingFileMetas(
partInfo.get(), bucket),
e);
}
if (restored.totalBuckets() != null && validateNumBuckets) {
checkNumBuckets(partInfo.get(), expectedTotalBuckets, restored.totalBuckets());
Integer restoredTotalBuckets = restored.totalBuckets();
if (restoredTotalBuckets != null
&& validateNumBuckets
&& expectedTotalBuckets != restoredTotalBuckets) {
if (partitionType.getFieldCount() > 0
&& options.bucketPerPartitionCountEnabled()
&& !strictBucketCount) {
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.",
Comment on lines +707 to +720

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.

partInfo.get(),
restoredTotalBuckets,
expectedTotalBuckets);
} else {
checkNumBuckets(partInfo.get(), expectedTotalBuckets, restoredTotalBuckets);
}
}
return restored;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import org.apache.paimon.mergetree.compact.CompactRewriterFactory;
import org.apache.paimon.metrics.MetricRegistry;
import org.apache.paimon.table.sink.CommitMessage;
import org.apache.paimon.table.sink.PartitionBucketMapping;
import org.apache.paimon.table.sink.SinkRecord;
import org.apache.paimon.types.RowType;
import org.apache.paimon.utils.CommitIncrement;
Expand All @@ -53,6 +54,12 @@ public interface FileStoreWrite<T> extends Restorable<List<FileStoreWrite.State<

FileStoreWrite<T> withWriteRestore(WriteRestore writeRestore);

/** Provides the preloaded partition-to-bucket mapping for fixed-bucket writes. */
default FileStoreWrite<T> withPartitionBucketMapping(
PartitionBucketMapping partitionBucketMapping) {
return this;
}

FileStoreWrite<T> withIOManager(IOManager ioManager);

/** Specified the write rowType. */
Expand Down
Loading
Loading