Repository navigation
Conversation
|
Production review at [P1] Keep compaction and recovery out of the ambiguous-routing guard
This is reachable through supported Flink paths: 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 [P2] Preserve the promised no-op for unpartitioned tables in Spark
I exercised actual SQL with 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. |
|
Additional [P1] Compute routing from the default-normalized row at The new mapped path calls 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 rowsThe null input routes to bucket 0, while the stored default 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. |
|
Reviewed [P1] Retire obsolete bucket IDs when restoring maintenance state ( I captured real checkpoint state through both sink implementations’ 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 ( 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 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: 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. |
|
Thanks @JingsongLi I updated PR :) |
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 = trueon a partitioned fixed-bucket table:totalBucketsvalue to the writer.INSERT OVERWRITEroutes with the target layout so a rescaled partition’s rewritten files are both hashed and stamped with the new bucket count.(partition, bucket)are rejected when the option is enabled, because they cannot prove which bucket count was used for routing.The option remains a no-op for unpartitioned tables, which retain the existing single table-level bucket-count invariant.
Documentation
bucket.per-partition-count-enabledto generated core configuration documentation.savepoint stop → rescale/overwrite → restart.
Tests
spark.paimon.write.use-v2-write=falseandtrue.Verification