Skip to content

feat(datafusion): support the $buckets system table - #934

Merged
JingsongLi merged 2 commits into
apache:mainfrom
jackylee-ch:feat/datafusion-buckets-system-table
Oct 2, 2026
Merged

JingsongLi merged 2 commits into
apache:mainfrom
jackylee-ch:feat/datafusion-buckets-system-table

Conversation

@jackylee-ch

Copy link
Copy Markdown
Contributor

Java exposes BucketsTable (<table>$buckets) for per-bucket file statistics — the standard way to diagnose bucket skew and decide whether a bucket needs compaction. The DataFusion integration had no equivalent, so that diagnosis was unreachable from DataFusion.

This adds the $buckets provider: it aggregates the table's data files by (partition, bucket) into record_count, file_size_in_bytes, file_count and last_update_time, mirroring Java's schema and partition/bucket ordering. Rows come from the same scan $files already uses; this only groups them, and reads fail closed under query-auth like the other metadata tables.

Tested by aggregating $files over the shared fixture and asserting $buckets reproduces it exactly.

@JingsongLi

Copy link
Copy Markdown
Contributor

Requirement fit: SUPPORTED; bucket skew diagnostics are useful. However, P1: group by the raw partition row, not its rendered string. collect_bucket_rows uses (Option<String>, i32) as the map key. The Java row formatter is not injective: partitions (p1='a, b', p2='c') and (p1='a', p2='b, c') both render {a, b, c}. I reproduced this through DataFusion SQL at 3551282f with a two-column partitioned table (bucket=1, bucket-key=id): inserting one row in each partition yields 2 $files rows but only 1 $buckets row. The new table silently combines unrelated partitions and reports incorrect counts/sizes. Please group by BinaryRow plus bucket, then render the partition only for output; add this two-partition regression.

The original $buckets end-to-end test, query-authorization fail-closed test, formatting, and diff checks pass. I restored the temporary reproducer after running it. Java BucketsTable / BucketEntry group by the binary partition before formatting: https://github.com/apache/paimon/blob/master/paimon-core/src/main/java/org/apache/paimon/manifest/BucketEntry.java

jackylee-ch added a commit to jackylee-ch/paimon-rust that referenced this pull request Sep 26, 2026
Review follow-up (apache#934): `collect_bucket_rows` keyed its aggregation on the
rendered partition string, but the Java cast-to-string formatter is not
injective — `(p1='a, b', p2='c')` and `(p1='a', p2='b, c')` both render
`{a, b, c}` — so distinct partitions merged into one `$buckets` row with
combined counts. Key on the serialized `BinaryRow` bytes instead, keeping the
string for output and ordering; extract `aggregate_bucket_rows` so the
look-alike case has a focused regression test.
Java exposes BucketsTable (`<table>$buckets`) for per-bucket file
statistics -- the standard way to diagnose bucket skew and decide whether
a bucket needs compaction. The Rust DataFusion integration had no
equivalent, so that diagnosis was unreachable from DataFusion.

Add the provider: aggregate the table's data files by (partition, bucket)
into record_count, file_size_in_bytes, file_count and last_update_time,
mirroring Java's schema and partition/bucket ordering. The rows come from
the same scan the $files table already uses; this only groups them. Reads
fail closed under query-auth like the other metadata tables.
Review follow-up (apache#934): `collect_bucket_rows` keyed its aggregation on the
rendered partition string, but the Java cast-to-string formatter is not
injective — `(p1='a, b', p2='c')` and `(p1='a', p2='b, c')` both render
`{a, b, c}` — so distinct partitions merged into one `$buckets` row with
combined counts. Key on the serialized `BinaryRow` bytes instead, keeping the
string for output and ordering; extract `aggregate_bucket_rows` so the
look-alike case has a focused regression test.
@jackylee-ch
jackylee-ch force-pushed the feat/datafusion-buckets-system-table branch from cd138ab to 7d5c3de Compare September 30, 2026 13:37
@jackylee-ch

Copy link
Copy Markdown
Contributor Author

Addressed, and rebased onto current main — the branch now sits on top of the merged $aggregation_fields table, and both system tables register side by side.

collect_bucket_rows no longer groups by the rendered partition string. aggregate_bucket_rows now keys on the serialized BinaryRow bytes (which are injective), keeping the rendered string only for output and ordering. So (p1='a, b', p2='c') and (p1='a', p2='b, c') stay separate even though both render {a, b, c}. This matches Java, which groups by the binary partition before formatting (as you noted for BucketEntry); the rows still sort by partition string then bucket, with the raw bytes as a deterministic tie-break for look-alikes.

Regression: look_alike_partitions_do_not_merge builds exactly those two partitions, asserts the precondition that their rendered strings collide (so string keying would merge them), and asserts aggregate_bucket_rows returns 2 rows. I verified it is non-vacuous: keying on the string alone makes it fail with rows.len() == 1; restoring the byte key passes. I placed the regression at the aggregation boundary rather than SQL so the collision is exercised deterministically, without depending on how a comma-bearing partition value round-trips through INSERT.

The $aggregation_fields merge only touched the shared registration list and the test-file tail, so the rebase was mechanical. system_tables (24 tests, including $aggregation_fields and $buckets) and the buckets unit tests pass; clippy -p paimon-datafusion --all-targets --features fulltext,vortex -D warnings is clean.

@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.

+1

@JingsongLi
JingsongLi merged commit 297965a into apache:main Oct 2, 2026
14 checks passed
jackylee-ch added a commit to jackylee-ch/paimon-rust that referenced this pull request Oct 2, 2026
Review follow-up (apache#934): `collect_bucket_rows` keyed its aggregation on the
rendered partition string, but the Java cast-to-string formatter is not
injective — `(p1='a, b', p2='c')` and `(p1='a', p2='b, c')` both render
`{a, b, c}` — so distinct partitions merged into one `$buckets` row with
combined counts. Key on the serialized `BinaryRow` bytes instead, keeping the
string for output and ordering; extract `aggregate_bucket_rows` so the
look-alike case has a focused regression test.
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.

2 participants