Skip to content

[DISCUSSION] Extending Partitioning to Support More Variants #21992

Description

@gene-bordegaray

I’d like to restart the partitioning side of this discussion separately from the dynamic-filter PRs and threads (see #21207 for more background)

The main issue seems to be that DataFusion cannot always represent the physical partitioning that some data sources actually have. In our case, range-partitioned data has had to claim Partitioning::Hash, which lets us avoid repartitioning but makes later optimizer and dynamic-filter decisions brittle due to this false advertisement.

Today we are in this scenario:

Data is actually range-partitioned:
- partition 0: key < 100
- partition 1: 100 <= key < 200
- partition 2: 200 <= key < 300

To mimic this we tell DataFusion our partitioning properties are:
- Partitioning::Hash([key], 3)

Those are not the same partitioning scheme.

The general goal should be that DataFusion should be able to describe physical partitioning truthfully, be able to compare two partitionings for compatibility, and use that information only when compatibility is proven.

This is well shown in dynamic filtering:

Build side partitioning:
  p0: key < 100
  p1: 100 <= key < 200
  p2: 200 <= key < 300

Probe side partitioning:
  p0: key < 100
  p1: 100 <= key < 200 
  p2: 200 <= key < 300

These are compatible as every partition X on the build side is compatible with partition X on the probe side. With this information, optimizations can safely use partition-local behavior:

if partitioning is compatible:
  use partition-local filter (more selective)
else:
  use global/safe fallback

A possible direction is to evolve Partitioning toward an extensible abstraction, such as a PhysicalPartitioning trait, and add range partitioning as the first use case. Longer term, built-ins like hash partitioning could move behind the same abstraction.

This would give us a cleaner foundation for:

Does this direction seem interesting to others in the community nd any thoughts on this proposal?

cc: @NGA-TRAN @alamb @jayshrivastava @gabotechs @adriangb

Activity

  1. NGA-TRAN commented on May 2, 2026

    @NGA-TRAN
    Contributor

    Great description. The exact use case at DataDog. We plan to use this a lot.

  2. gene-bordegaray commented on May 3, 2026

    @gene-bordegaray
    ContributorAuthor

    I have added a draft PR #22002 which shows my vision for what partitioning can look like going forward and how this will apply to an optimization like dynamic filter routing. The API design was done by me but implementation details by AI and is not meant to be a mergeable PR. I provided a description of the API, would love any suggestions/ideas on this general approach.

    I separated it into 4 commits to show the logical progression of the implementation:

    1. Adding Range partitioning variant that models what a custom partitioning trait would look like
    2. Introduce the general partitioning trait using this model and have range partitioning implement it
    3. Have file preserved partitioning use range partitioning rather than false advertising hash
    4. Route dynamic filters using partitioning compatibility and create partitioning-local filters when range partitioned
  3. alamb commented on May 4, 2026

    @alamb
    Contributor

    I left some comments on #22002

  4. alamb commented on May 4, 2026

    @alamb
    Contributor

    I think it would also help to decide sooner rather than later if we want a general partitioning trait, or if we are just going to hard code Range Partitioning. I think either could work.

    If we are going to go with a trait, it might be good to declare that we eventually want to shoot to remove all special cases for hash partitioning 🤔

  5. NGA-TRAN commented on May 4, 2026

    @NGA-TRAN
    Contributor

    I think it would also help to decide sooner rather than later if we want a general partitioning trait, or if we are just going to hard code Range Partitioning. I think either could work.

    I vote for a general partitioning trait

  6. adriangb commented on May 4, 2026

    @adriangb
    Contributor

    I'm not sure if it should be a trait or an eum. I generally prefer an explicit enum if the trait doesn't generalize well. E.g. if we can't point to an obvious 3rd use case for the trait that an external system could conceivably implement and where everything would work for them, I fear a trait is the wrong choice.

    @NGA-TRAN could you share why you prefer a trait to an enum here? Maybe you know some additional partitioning strategies that are commonly used in the database world that I don't know of.

  7. NGA-TRAN commented on May 4, 2026

    @NGA-TRAN
    Contributor

    @adriangb : Can you provide an example of enum? If it covers our use cases, we are happy to go with it.

    Currently, our range come from 2 columns (or 2 expressions). For example (time, int_val). So we can have these ranges:

    • ( [2026.05.01 - 2026.05.07] , [1 - 100] )
    • ( [2026.05.01 - 2026.05.07] , [101 - 200] )
    • ...
    • ( [2026.05.08 - 2026.05.15] , [1 - 500] )
    • ( [2026.05.08 - 2026.05.15] , [501 - 1000] )
    • ...

    As long as they do not overlaps, they are considered as ranges. We can definitely map those into simple ranges or indexes and use those simple ones in DF instead.

  8. adriangb commented on May 4, 2026

    @adriangb
    Contributor

    I don't know that the enum is any different than the trait in terms of representation. I.e. if we can represent it as a trait we can represent it as an enum. The question is more if we want to hard code all of the handled cases or leave it open for users to add more down the road.

    I'd defer to @gene-bordegaray as to what an enum version would look like but I think something like this:

      pub enum Partitioning {
          RoundRobinBatch(usize),
          Hash {
              exprs: Vec<Arc<dyn PhysicalExpr>>,
              partition_count: usize,
          },
          Range {
              exprs: Vec<Arc<dyn PhysicalExpr>>,
              partitions: Vec<RangePartition>,
          },
      }
    
      pub struct RangePartition {
          /// One interval per `Range.exprs` entry.
          ///
          /// For exprs = [time, int_val], this can represent:
          /// time in [2026-05-01, 2026-05-07]
          /// AND int_val in [1, 100]
          pub bounds: Vec<RangeInterval>,
      }
    
      pub struct RangeInterval {
          pub lower: Option<RangeBound>,
          pub upper: Option<RangeBound>,
      }
    
      pub struct RangeBound {
          pub value: ScalarValue,
          pub inclusive: bool,
      }
  9. NGA-TRAN commented on May 4, 2026

    @NGA-TRAN
    Contributor

    The enum Gene proposed looks good. As long as there are PhysicalExpr, it will work

    Range {
              exprs: Vec<Arc<dyn PhysicalExpr>>,
              partitions: Vec<RangePartition>,
          },
  10. gene-bordegaray commented on May 5, 2026

    @gene-bordegaray
    ContributorAuthor

    I am in favor of a general trait as the long-term goal of this work. I think allowing users to implement their own type of partitioning will make DF more powerful in production use cases as I am sure that there will be instances of partitioning that are not captured in Hash or Range partitioning (just as Hash did not fully work for us). Off the top of my head something like a Value partitioning would also be useful:

    p0: col in ('a', 'd')
    p1: col in ('b')
    p2: col in ('c')
    

    Because of this I think providing another extendible point for people (the trait) will be very high value even if worth the extra effort.

    With this said I do think we can create mergeable commits by extending the enum now by supporting Range partitioning as @adriangb has described but model it after what our trait will look like. We can treat the trait as the final goal but let Range help us define the requirements for that as we pseudo-implement what that trait will look like.

    If we are going to go with a trait, it might be good to declare that we eventually want to shoot to remove all special cases for hash partitioning 🤔

    And along with this, yes I agree here that we should shoot to encapsulate all partitioning logic behind this trait and no special cases. The optimizer rules and other things should ask if two partitioning are compatible or satisfy one another, not just "is this Hash partitioned" 👍

    I see this path as actually being more intuitive once done well.

  11. stuhood commented on May 12, 2026

    @stuhood
    Contributor

    I think that partitioning is likely to become relevant for us in the next few weeks, and our schedule should finally be clearing up to try and get some of our team involved with helping.

    We're very likely to be doing either Range or N-dimensional partitioning. AFAICT though, Range is sufficient to represent N-dimensional partitioning (with a physical optimizer rule running before EnforceDistribution), because you can declare each table to be partitioned on the relevant 1-dimensional join key, as a subset of the N-dimensions.

    So, in practice, it doesn't seem like it would be worth actually introducing knowledge of multi-dimensional ranges here: a Range type would be enough.

  12. alamb commented on May 14, 2026

    @alamb
    Contributor

    We spoke about this topic the other day in person when we were at the NYC meetup

    I believe @gene-bordegaray 's usecase will be satisfied if we can represent range partitioning that is represented like

    • [expr]: [[ranges]]

    For example, partitioning on date would be represented like

    • [date]: [
      [[2021-01-01, 2021-12-32]]
      [[2022-01-01, 2022-12-32]],
      ...
      ]

    partitioning on date, city would be represented like

    • [date, city]: [
      [[2021-01-01, 2021-12-32], [Allston, Boston]]
      [[2021-01-01, 2021-12-32], [Boston, NYC]]
      [2022-01-01, 2022-12-32], [Allston, Boston]]
      [2022-01-01, 2022-12-32], [Boston, NYC]]
      ...
      ]

    Some TBD details are:

    1. How to ensure the entire range is covered (e.g. what range in the above example handles 1970-01-01)?
  13. alamb commented on May 14, 2026

    @alamb
    Contributor

    (would the above work for you and your usecse @stuhood ?)

  14. stuhood commented on May 14, 2026

    @stuhood
    Contributor

    In my experience, range partitioning is usually represented with only the partitioning points, which addresses the infinite end cap issue, and guarantees that everything is covered. Each point is exclusive for the range to the left side of the point and inclusive for the range to the right side of the point. So your example would be more like:

    To create three ranges (infinitely small to 2021-01-01, between 2021-01-01 and 2022-12-32, and 2022-12-32 to infinitely large)

    `[date]`: [
      [2021-01-01]
      [2022-12-32], 
      ...
      ]
    

    And the date + city example would be:

    `[date, city]`: [
      [2021-01-01, Allston],
      [2021-01-01, Boston],
      [2022-12-32, Allston],
      [2022-12-32, NYC],
      ...
      ]
    

    (or something).

  15. gene-bordegaray commented on May 15, 2026

    @gene-bordegaray
    ContributorAuthor

    To keep everyone else in this thread up-to-date, we had some discussions this week regarding this topic and have come up with a concrete plan. We are going to start with adding a enum variant for ExprPartitioning (name up for debate) rather than introducing a general trait immediately. @alamb described the concept and representation of this well, so I won't reexplain that here.

    For implementation, the first PR will be purely mechanical. This will add a Expr (or similar) variant to the existing physical Partitioning enum, along with the supporting types, but just throw "not implemented" at callsites. This will introduce the concept and give a good idea of the API methods we will need to represent partitioning well.

    This will then be followed up with implementing the partitioning features based on the callsites in follow-up PRs, with one of these beeing the dynamic filter routing.

    The goal is still to model this in a way that can evolve into a general partitioning abstraction if needed.

    I will make the mechanical PR soon and ping here 👍

  16. 9 remaining items

  17. alamb commented on Aug 25, 2026

    @alamb
    Contributor

    Given #22395 maybe we should close this issue and track the work there

  18. gene-bordegaray commented on Aug 25, 2026

    @gene-bordegaray
    ContributorAuthor

    Given #22395 maybe we should close this issue and track the work there

    sounds good to me

  19. alamb commented on Aug 25, 2026

    @alamb
    Contributor

    Tracking in #22395

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions