Skip to content

Enable parquet filter pushdown (filter_pushdown) by default #3463

Description

@alamb

In #3380 @thinkharderdev added support for evaluating filters during the parquet scan via the RowIndex mechanism 🎉

This feature is currently enabled via a feature flag, which is disabled by default.

This ticket tracks enabling this feature by default.

Currently known items are:

Activity

  1. self-assigned this
    on Oct 13, 2022
  2. alamb commented on Oct 13, 2022

    @alamb
    ContributorAuthor

    I plan to review the parquet test coverage over the next day or two

    My basic plan is to:

    1. Make a draft PR that enables pushdown by default and see if any regression tests fail
    2. Do some manual testing using datafusion-cli and some internal datasets
  3. alamb commented on Oct 13, 2022

    @alamb
    ContributorAuthor

    The PR that enables the feature is #3828 and it looks quite promising

    Also, I have a proposed config change that I would like to get in, #3822, which would allow users to quickly disable this feature if they found issues

  4. alamb commented on Oct 26, 2022

    @alamb
    ContributorAuthor

    Update here: I am working on some fuzz testing for the predicate pushdown code

  5. alamb commented on Oct 28, 2022

    @alamb
    ContributorAuthor

    I found a bug in pushdown by writing a test #4005

  6. alamb commented on Oct 28, 2022

    @alamb
    ContributorAuthor

    I found another one with my test. I will keep the list on this ticket updated

  7. alamb commented on Oct 31, 2022

    @alamb
    ContributorAuthor

    I think we are getting close

  8. alamb commented on Nov 2, 2022

    @alamb
    ContributorAuthor

    Update here is I think once we get #3976 merged I'll put up the PR to enable the feature by default

  9. Dandandan commented on Feb 5, 2023

    @Dandandan
    Contributor

    @alamb What's the remaining work that needs to be done to enable it by default? Anything I can do to get this over the finish line?

  10. alamb commented on Feb 5, 2023

    @alamb
    ContributorAuthor

    @Dandandan ❤️ thank you

    I know of two major items:

    1. Enable page index by default Enable parquet page level skipping (page index pruning) by default #4085
    2. Ensure that enabling filter pushdown by default doesn't cause performance regressions in cases where the filter isn't very selective. I know @tustvold had some thoughts in this area.

    For item 1, I tried a few times, but am currently blocked by #5104 and haven't had time to return to it

    For item 2, I think someone needs to look at some benchmark results (perhaps TPCH) and figure out what we can do to avoid the regression (or measure and determine it isn't significant)

    I keep hoping to have time to work on item 1, as I am pretty sure I know what the problem is, but other things keep coming up :(

    I haven't had a chance to work on 2 yet

  11. alamb commented on Mar 2, 2023

    @alamb
    ContributorAuthor

    Update here is I am working on a larger benchmarking story, part of which would give us more confidence to merge changes like this in. I hope to have that done early next week

  12. alamb commented on Jul 24, 2023

    @alamb
    ContributorAuthor

    With the introduction of bench.sh with tpch and clickbench (#6994) I think we are now in a better position to do this experiment again and fix any performance regressions

  13. tustvold commented on Mar 17, 2024

    @tustvold
    Contributor

    apache/arrow-rs#5523 might help mitigate the impact of pushing down predicates that turn out to not be very selective

  14. alamb commented on Mar 17, 2024

    @alamb
    ContributorAuthor

    It would be really nice to (finally) be able to turn this optimization on -- thank you @tustvold

    cc @Dandandan

  15. removed their assignment
    on Aug 5, 2024
  16. 23 remaining items

  17. tustvold commented on Jan 4, 2026

    @tustvold
    Contributor

    Yeah, it's a good point that whilst caching reduces the additional decode costs for pushing down predicates, it doesn't eliminate the IO costs. That being said in general you only really want to be pushing down one or two ArrowPredicate as each has costs associated with it, with this just becoming even more true against object stores.

    One has 3 choices with what to do with a given predicate:

    1. Push it down as-is as an ArrowPredicate
    2. Fuse it with some other predicates into a combined ArrowPredicate
    3. Don't push the predicate down at all

    If you have multiple ArrowPredicate you then have the additional complexity of what order to apply them in.

    I'm not familiar with how DF is handling this currently, but a selectivity estimate based approach at plan time might be a good place to start.

  18. adriangb commented on Jan 4, 2026

    @adriangb
    Contributor

    I'm not familiar with how DF is handling this currently, but a selectivity estimate based approach at plan time might be a good place to start.

    The answer is: we are not. The only similar thing we do is use the column sizes (from parquet metadata) to reorder the filters. Otherwise we split the conjunction (split at and operators) and each clause gets evaluated independently, ordered from smallest column size to largest.

    I don’t think we have enough information to do anything useful from statistics (this is probably why we haven’t done so yet) but if arrow-rs at least exposed the selectivity of filters after each file is read (ideally each batch?) we could at least have runtime filter selectivity statistics so as we open more files we adapt our approach using the options you described above. A further step would be for arrow-rs to allow us to rebuild/reshuffle our approach within a scan but that may require more API churn. Adjusting between files should be pretty simple.

  19. tustvold commented on Jan 4, 2026

    @tustvold
    Contributor

    arrow-rs at least exposed the selectivity of filters after each file is read

    It is possible to provide an implementation of ArrowPredicate that tracks this. IIRC there even is a test in the parquet crate that does just this.

    Adjusting between files should be pretty simple.

    This sounds like a nice idea, and I agree should be a relatively straightforward lift.

    I don’t think we have enough information to do anything useful from statistics

    I'm no expert here, but this seems off to me. Parquet provides metadata about sort orders, min/max values, null counts, etc...

    Unfortunately distinct_count is rarely populated, but I think you should be able to do something by checking to see if a column spilled its dictionary. If it did either the values are big or there are lots of them - both making this a poor pushdown cancidate. If it didn't spill its dictionary, this might be a good enough signal on its own, but you may be able to inspect the offset index (I can't remember if it includes dictionary pages) to get the precise NDV.

    To say there is not enough information, seems overly pessimistic... You can always refine your initial guess later as you describe.

  20. adriangb commented on Jan 4, 2026

    @adriangb
    Contributor

    Okay yes I agree maybe I was being pessimistic 😆. In any case using what we can from stats / metadata to set up the initial state / plan and then refining it once we have runtime statistics seems like a good place to land, and those efforts can be tackled orthogonally.

  21. adriangb commented on Jan 4, 2026

    @adriangb
    Contributor

    arrow-rs at least exposed the selectivity of filters after each file is read

    It is possible to provide an implementation of ArrowPredicate that tracks this. IIRC there even is a test in the parquet crate that does just this.

    Any chance you could link to that test? Would be sweet if it’s doable with the current API!

  22. tustvold commented on Jan 4, 2026

    @tustvold
    Contributor

    Any chance you could link to that test? Would be sweet if it’s doable with the current API!

    I can't find it with a quick scan, but you can definitely do it using some sort of shared state that you increment in the predicate closure.

    For example

    struct State {
        filtered: AtomicUsize,
        passed: AtomicUsize,
    }
    
    let state = Arc::new(State::default());
    
    let s_captured = Arc::clone(&state);
    let filter = Box::new(ArrowPredicateFn::new(projection, move |batch|{
        let filtered = filter(batch);
        state.filtered.fetch_add(Ordering::Relaxed, batch.len());
        state.passed.fetch_add(Ordering::Relaxed, filtered.len());
        return Ok(batch)
    }))
    

    You could even wrap this up in a custom impl ArrowPredicate instead of relying on ArrowPredicateFn.

  23. alamb commented on Jan 5, 2026

    @alamb
    ContributorAuthor

    Oh right, yes it will do that sorry, been years since I wrote that code (and it looks like there's some new PushDecoder anyway that might change all of this).

    FWIW the push decoder doesn't change the fundamental IO -- it is instead an explicit state machine (rather than using the implicit one created by the rust compiler with await calls). It is algorithmically the same thing

    So yes it will behave the way you describe, which could make a difference for stores with high first-byte latencies. However, my understanding is Andrew is running on an NVMe drive where time spent doing IO will be bounded by bytes-read not number of fetches, which filter pushdown would if anything reduce...

    The idea that one reason it slows down is due to additional latency of multiple distinct IOs is an interesting one. I will ponder that and see if I can find some way to reduce it / prefetch

  24. adriangb commented on Jan 21, 2026

    @adriangb
    Contributor

    @alamb just to unblock you here if dynamic filter pushdown from hash joins is the cause of the slowdown it would be reasonable to disable it by default

  25. alamb commented on Jan 22, 2026

    @alamb
    ContributorAuthor

    (BTW @Dandandan @lyang24 and others are currently "optimizing the s#!t" out of the Parquet reader so we'll see how things look with arrow 58)

  26. Dandandan commented on Feb 3, 2026

    @Dandandan
    Contributor

    I think I have mostly traced down the slowdown of TPCH and (dynamic) filter pushdown:

    • Dynamic filter pushdown creates a long case when expression based on the number of target partitions (~75% or more of the overhead), where the evaluation cost grows with the number of partitions.
      Evaluating this expression could be optimized to use direct indexing (instead of filtering per expression) to make overhead really small (I'll be creating a PR later this week for this, currently fighting some illness).
    • contains_hashes is relatively expensive if nothing can be filtered out. I think we should disable the hash table push down by default for now (only enable the direct indexing ArrayMap if no slowdown happens there or disable it altogether for the moment) and implement some better heuristics to enable / disable the hash table lookup (based on some sampling?)
  27. Dandandan commented on Feb 3, 2026

    @Dandandan
    Contributor

    Thinking about it a bit more, I think the fastest way forward is to disable the evaluation of the dynamic filter predicate pushdown in the scan / predicate pushdown:

    • The evaluation of the big expression case hashes % partitions when 0 then (...) AND lookup when 1 ... is expensive, even with optimizations, as also the branching expressions are different (based on the min/max values in the hash maps) -> I think we can reduce this by comparing against the global/combined statistics rather than a per-partition statistic and implementing a fast way to check a batch of hashes against n tables
    • The expression and contains_hashes will always have overhead, so we need first a smart way to dynamically disable it when it doesn't filter out anything / a lot.
  28. adriangb commented on Feb 3, 2026

    @adriangb
    Contributor

    I'm +1 on disabling the hash join pushdown by default if it's slow.

    I think we can then split work up into 3 streams:

    • Rework the expressions themselves.
      • Pull statistics out of the hash / case expressions.
        • Global statistics guard: col < 123 AND col > 5 AND CASE ... END
        • Per-partition statistics guard: CASE WHEN col < 23 AND COL > 15 AND hash(col) % partitions = 0 THEN lookup ....
    • Optimize expression evaluation (make hash computation more efficient, optimize the CASE expression, etc.)
    • Optimize expression optimization (the work you've been doing on PhysicalExprSimplifier)
    • Disable the expressions dynamically.

    For that last point I think it should be pretty easy to code up a wrapper PhysicalExpr:

    struct Selectivity {
        rows_processed: usize,
        rows_matched: usize,
    }
    
    struct OptionalFilterPhysicalExpr {
       target_selectivity: f64,
       inner: RwLock<Option<Arc<dyn PhysicalExpr>>>,
       selectivity: RwLock<Selectivity>,
    }
    
    impl PhysicalExpr for OptionalFilterPhysicalExpr {
       fn evaluate(&self, batch: RecordBatch) -> Result<RecordBatch> {
          // if inner is None return all true
          // otherwise evaluate with inner and update selectivity
         // if selectivity is < target after evaluating X rows make inner None
         // (which also drops the reference so that any producers of dynamic filters can stop updating them if they want)
      }
    }
  29. alamb commented on Feb 12, 2026

    @alamb
    ContributorAuthor
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