Skip to content

Further improve performance of IN list evaluation #19241

Description

@geoffreyclaude

Summary

IN LIST evaluates expressions such as:

x IN (1, 3, 7)

When the right-hand list is constant, DataFusion can build a membership filter once and reuse it for every input row. Dynamic-filter pushdown can apply the same check millions of times during a scan.

This epic adds exact lookup strategies selected by physical representation and non-null list length. If none applies, DataFusion uses the generic static filter.

All strategies share the same result construction, so this work does not change SQL behavior for IN, NOT IN, input nulls, nulls in the list, dictionaries, or sliced arrays.

Stack

The PRs should be read in this order.

Landed

Remaining

How the strategies work

Generic fallback

#21927 improves the path available to every supported type. It precomputes Arrow hashes for the constant list, stores list indexes in a compact hash table, and uses Arrow's exact comparator to confirm equality. Membership is produced as a bitmap and then combined with input validity, list nulls, and NOT IN semantics by shared result-building code.

Bitmap lookup for small fixed domains

A one- or two-byte value has only 256 or 65,536 possible bit patterns. That entire domain can be represented by one bit per pattern:

  • UInt8 and Int8 use a 256-bit (32-byte) bitmap.
  • UInt16, Int16, and Float16 use a 65,536-bit (8 KiB) bitmap.

Building the filter sets the bit corresponding to each non-null list value. Probing a row is then one indexed bit test, with no hashing and no scan of the list. Signed integers and Float16 map their native bit patterns directly into the same finite bitmap domain.

Direct comparisons for very small primitive lists

For a tiny fixed-width list, a hash lookup can cost more than comparing the input directly with every constant. #23014 stores the constants in a fixed-size array and combines all equality checks into a predictable comparison chain.

The maximum non-null list sizes are:

Physical width Direct comparisons through
1 byte 16 values
2 bytes 8 values
4 bytes 32 values
8 bytes 16 values
16 bytes 4 values

Larger one- and two-byte lists use the bitmap strategy. Larger native 32- and 64-bit integer lists use primitive hash sets; #24283 extends that same shared fallback to Decimal128. Other unsupported primitive types continue through their existing exact filter or the generic fallback.

Shared primitive selection and Decimal128

#24283 puts the direct-comparison thresholds and larger-list choices in one primitive selector. Native arrays and representation adapters call that selector, so there is no second dispatch table.

For native Decimal128, up to four non-null list values use direct comparisons. Larger lists store and probe the native unscaled i128 values in the shared primitive hash-set filter instead of using the generic Arrow array filter. Only membership storage changes; decimal expression typing, precision/scale compatibility, null behavior, and exact native-value equality remain unchanged.

FixedSizeBinary

#24102 recognizes FixedSizeBinary widths that exactly match an existing primitive representation. The list and input are converted in the same way, so primitive equality and hashing are exact equality over the original fixed-width bytes; no numeric or decimal operations are involved.

FixedSizeBinary width Internal key Small lists Larger lists
1 byte UInt8 direct comparisons through 16 values bitmap
2 bytes UInt16 direct comparisons through 8 values bitmap
4 bytes UInt32 direct comparisons through 32 values shared primitive hash-set filter
8 bytes UInt64 direct comparisons through 16 values shared primitive hash-set filter
16 bytes 128-bit primitive direct comparisons through 4 values shared primitive hash-set filter

Other widths keep using the generic filter.

Aligned Arrow buffers are reused directly. If a valid array has an unaligned buffer, the adapter copies its values into aligned primitive storage first. Alignment affects whether a copy is needed, not correctness.

Inline Utf8View and BinaryView

Arrow stores a Utf8View or BinaryView value of at most 12 bytes completely inside its 128-bit view: the view contains the length and the zero-padded bytes. For these values, the view holds the full value, so view equality is exact and needs no backing-buffer read.

#24088 uses this representation only when every non-null value in the constant list is inline:

  • up to four non-null values use the existing 128-bit direct-comparison filter;
  • larger lists reuse the Decimal128 hash-set filter over the same 128-bit key; and
  • if any list value is longer than 12 bytes, the whole list uses the generic filter.

Input arrays may still contain long values. Their encoded length distinguishes them from every inline list key, so they produce a miss without reading their backing bytes. Utf8View and BinaryView remain separate typed paths, and dictionary inputs continue through the shared dictionary handling.

Expected impact

Case Lookup used
Generic constant lists precomputed Arrow hash table with exact comparison
One- and two-byte primitive domains one bitmap bit test per row
Very small primitive lists fixed direct-comparison chain
Larger native Decimal128 lists shared primitive hash-set lookup
Supported FixedSizeBinary widths reused primitive direct, bitmap, or hash-set lookup
All-inline Utf8View / BinaryView lists exact 128-bit direct or hash-set lookup

Activity

  1. adriangb commented on Dec 9, 2025

    @adriangb
    Contributor

    Amazing plan!

  2. Dandandan commented on Dec 9, 2025

    @Dandandan
    Contributor

    What about slice::contains? Seems like it should be somewhere between the const-sized approach and binary search in terms of threshold window.

  3. adriangb commented on Dec 9, 2025

    @adriangb
    Contributor

    There is a related bit of work to untangle with is the ScalarValue references that InListExpr is forced to use even if we are starting from an array. It uses these to look up in bloom filters, do predicate pruning, etc. We could make all of the relevant APIs work with an enum of Vec<ScalarValue> (heterogenous lists) or ArrayRefs (homogenous lists) and that would avoid converting array -> ScalarValue and then in some places back to an array (deep in pruning code iirc). Some of this is planning time / build time stuff so it is amortized over the data scans, but some of it happens for each file opened. It's not as big as for each row, but it adds up.

    A second thing is that col IN (...) is inefficient when it hits a bloom filter on col if the list is large because it loops over the values in the list. I'm not sure how or where we would do this but in theory we could build a bloom filter out of the InListExpr and then do a binary operation between that bloom filter and Parque's bloom filter instead of looping over each item in the list and looking it up in the columns bloom filter. At the very least if we did the point above and pushed down an array we could probably be more efficient about converting all of the values into something we can look up in the bloom filter (currently it goes through ScalarValue).

  4. geoffreyclaude commented on Dec 9, 2025

    @geoffreyclaude
    ContributorAuthor

    What about slice::contains? Seems like it should be somewhere between the const-sized approach and binary search in terms of threshold window.

    It loses all the time against the branchless and the binary search. slice::contains is just a wrapper around:

    self.iter().any(|y| y == x)

    So it needs to do bounds checks and can't optimize for the small arrays.

  5. Dandandan commented on Dec 9, 2025

    @Dandandan
    Contributor

    What about slice::contains? Seems like it should be somewhere between the const-sized approach and binary search in terms of threshold window.

    It loses all the time against the branchless and the binary search. slice::contains is just a wrapper around:

    self.iter().any(|y| y == x)
    So it needs to do bounds checks and can't optimize for the small arrays.

    Hmm interesting. I would have thought it would be faster somwehere in between.
    Note it is not equal to self.iter().any(|y| y == x), slice::contains is specialized for primitive types to create unrolled/vectorized code.

  6. geoffreyclaude commented on Dec 10, 2025

    @geoffreyclaude
    ContributorAuthor

    @Dandandan See geoffreyclaude#14 for an in-depth micro benchmark and analysis of the different search algorithms.

    TL;DR: It's always branchless up to the SIMD limit, then hashset.

    Slice Search Benchmark Image
  7. geoffreyclaude commented on Dec 17, 2025

    @geoffreyclaude
    ContributorAuthor

    I've opened #19376 as a preliminary PR to extend the benchmarks.

  8. alamb commented on Dec 18, 2025

    @alamb
    Contributor
    1. Type Normalization

    Note you can potentially use the RowFormat for this purpose - https://docs.rs/arrow-row/latest/arrow_row/, perhaps as a fallback when more general methods aren't available

    It handles pretty much every arrow type

    The downside is that the input needs to be first converted into row format

  9. 2 remaining items

  10. alamb commented on Sep 30, 2026

    @alamb
    Contributor
  11. geoffreyclaude commented on Sep 30, 2026

    @geoffreyclaude
    ContributorAuthor

    @geoffreyclaude -- is this epic done? It looks like we have these left?

    Yes, we still have these two. I need to update #25186 following reviews, and #25187 should follow. No reason they can't make it to the DF56 release!

  12. geoffreyclaude commented on Oct 9, 2026

    @geoffreyclaude
    ContributorAuthor

    @alamb Sorry for the delay, this was trickier than anticipated!

    I gave it some thought and suggest inverting the two PRs order. First the optimization OR-rewrite, and then the signed zero bug fix. Simply because the signed zero issue is pre-existing and might need a deeper discussion.

    I've rebased #25187 and opened it for review. It's not trivial though as the gains depend on the data type: for some we should keep the OR-rewrite.

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