Skip to content

feat(datafusion): bucket FileScanTasks across target_partitions with … - #1

Merged
toutane merged 3 commits into
toutane:draft/partitioned-file-scanning-contributionfrom
timsaucer:feat/partitioned-scan-hash-bucketing
Apr 28, 2026
Merged

toutane merged 3 commits into
toutane:draft/partitioned-file-scanning-contributionfrom
timsaucer:feat/partitioned-scan-hash-bucketing

Conversation

@timsaucer

Copy link
Copy Markdown

Which issue does this PR close?

None

What changes are included in this PR?

  • When possible set output partitioning to Hash based on the partitions.
  • Output the target_partitions or number of file scans, whichever is greater
  • Create even partitioning based on hash of file path when we cannot determine partitions so that we at least do better than round robin

Are these changes tested?

Some unit tests are included, but I need help to test against real world catalogs.

…identity-hash partitioning

Replace the one-task-per-partition layout in IcebergPartitionedScan with
N buckets sized from the session's target_partitions. When the table's
default spec exposes identity-transform columns and every task carries
the corresponding partition values, tasks are bucketed by hashing those
values via DataFusion's REPARTITION_RANDOM_STATE so the resulting
partitioning matches what RepartitionExec would produce. The scan then
declares Partitioning::Hash(exprs, N), letting downstream joins and
aggregates skip an extra repartition.

Hash declaration is conservative and only stands when:
  - the table has a single partition spec (no spec evolution)
  - every identity source column is present in the output projection
  - every column type is supported by literal_to_array
  - every task supplied a full identity key
Any miss collapses to UnknownPartitioning(N) while bucketing falls
back to a hash of data_file_path so partitions still distribute.

IcebergPartitionedScan now stores Vec<Vec<FileScanTask>> and execute(i)
streams every task in buckets[i] through to_arrow_with_tasks. Bucket
count is capped at min(target_partitions, num_files), and an empty
table still yields zero partitions to avoid out-of-bounds execute calls.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
timsaucer and others added 2 commits April 24, 2026 19:37
`IcebergPartitionedTableProvider::supports_filters_pushdown` previously
returned `Inexact` for every filter, forcing DataFusion to re-evaluate
even filters that Iceberg's manifest-level pruning has fully resolved.
Per-filter the provider now returns `Exact` when both:
  - the iceberg conversion can represent the filter, so manifest pruning
    will remove every row that fails it, and
  - every leaf is a comparison or null check against an identity-
    partition column with a literal RHS.

Identity-partitioned column names are cached at `try_new` from the
table's default spec; tables with spec evolution (>1 historical specs)
fall back to an empty set so all filters stay `Inexact`. Supported
shapes: =, !=, <, <=, >, >=, IS NULL, IS NOT NULL, IN/NOT IN, plus
AND/OR/NOT compositions of the above. Every other shape is `Inexact`.

`convert_filter_to_predicate` is promoted to `pub(crate)` so the
provider can probe convertibility per filter without rebuilding the
whole AND-collapsed predicate.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
…column intersection

Previously identity_partition_col_names returned an empty set whenever
the table had more than one historical partition spec, forcing every
filter back to Inexact under spec evolution. This was overly
conservative: Iceberg evaluates partition predicates against each
manifest's own spec, so a column that is identity-partitioned in every
spec is fully prunable across the entire table regardless of which spec
a given file was written under.

Replace the multi-spec gate with an intersection across every spec's
identity-source set. A column survives only if every spec includes it
with Transform::Identity; columns that appear with non-identity
transforms in some spec, or are missing from a spec entirely, are
dropped. The result remains an honest set of columns for which Exact
pushdown is provably safe across all surviving files.

Hash bucketing (compute_identity_cols) keeps its single-spec gate
because slot-order alignment with the table's default spec depends on
each task carrying its own spec id, which the native plan flow does
not yet do.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

@toutane toutane left a comment •

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

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

Thanks @timsaucer for this PR !

It is really clean and definitely replace the old table provider.

As for the filter pushdown optimization, maybe it would be a good idea to do that in a second pull request upstream, to keep the scope of this one limited? What do you think?


use super::*;

async fn make_catalog_and_partitioned_table(

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

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

Maybe it is worth replacing calls to make_catalog_and_table with make_catalog_and_partitioned_table(None) to remove ~40 lines of duplication.

@timsaucer

Copy link
Copy Markdown
Author

As for the filter pushdown optimization, maybe it would be a good idea to do that in a second pull request upstream, to keep the scope of this one limited? What do you think?

Sounds good, but if that's the way you want to go then I suggest opening a tracking issue so we don't lose it.

@toutane

toutane commented Apr 27, 2026

Copy link
Copy Markdown
Owner

As for the filter pushdown optimization, maybe it would be a good idea to do that in a second pull request upstream, to keep the scope of this one limited? What do you think?

Sounds good, but if that's the way you want to go then I suggest opening a tracking issue so we don't lose it.

@timsaucer, I created this ticket last Friday: apache#2363

I was suggesting a second pull request because I think it would help us merge the partitioned scan more quickly? But I might be wrong; you’re much more familiar with this kind of process, so let me know if you don’t think a follow-up PR is worth it.

@toutane
toutane merged commit 6d8888f into toutane:draft/partitioned-file-scanning-contribution Apr 28, 2026
17 checks passed
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