Skip to content

Remove the IN predicate limit so large IN predicates prune files - #4032

Open
azwanzuharimi wants to merge 2 commits into
apache:mainfrom
azwanzuharimi:fix-in-predicate-pruning
Open

azwanzuharimi wants to merge 2 commits into
apache:mainfrom
azwanzuharimi:fix-in-predicate-pruning

Conversation

@azwanzuharimi

@azwanzuharimi azwanzuharimi commented Sep 30, 2026 •

Copy link
Copy Markdown

Closes #4003

Rationale for this change

_InclusiveMetricsEvaluationVisitor.visit_in and _ManifestEvalVisitor.visit_in return ROWS_MIGHT_MATCH as soon as an In predicate has more than IN_PREDICATE_LIMIT (200) values. They return before they read any bounds. So plan_files keeps every data file above that limit. A delete or an upsert with a large key set then scans the whole table. #4003 measured this: 200 keys plan 1 file out of 20, 201 keys plan all 20.

The limit copies Java InclusiveMetricsEvaluator, where apache/iceberg#1672 added it in 2020. The concern was the cost of filtering every literal per file.

This change removes the two early returns. Both visitors now run the existing exact check for all sizes. The cost is one scan of the value set per file. IN_PREDICATE_LIMIT stays in the module because the tests use it, but no code reads it now.

Measured on a local table with 20 data files of 100,000 rows each. The keys all come from one file. Planning time is near equal. The scan reads fewer files.

keys files planned before files planned after scan time before scan time after
201 20 / 20 1 / 20 22 ms 12 ms
1,000 20 / 20 1 / 20 39 ms 14 ms
10,000 20 / 20 1 / 20 239 ms 37 ms
100,000 20 / 20 1 / 20 2,234 ms 279 ms
keys plan before plan after
201 8.4 ms 8.9 ms
10,000 23 ms 25 ms
100,000 154 ms 171 ms

table/update/validate.py uses the same evaluator for conflict detection. It drops a file only on ROWS_CANNOT_MATCH, so this change can only remove false conflicts.

This diverges from Java, which still returns ROWS_MIGHT_MATCH above the limit.

Are these changes tested?

Yes.

  • tests/expressions/test_evaluator.py: test_integer_in_above_limit, test_integer_in_at_limit and test_inclusive_metrics_evaluator_in_above_limit. They cover values below, above, equal to, across and outside both bounds, the 200 value boundary, an all nulls column, and NaN bounds.
  • tests/expressions/test_visitors.py: test_integer_in_above_limit for the manifest evaluator.

The above limit tests fail on main. make lint passes. tests/expressions and tests/table/test_evaluator_planning.py pass with 588 tests.

Are there any user-facing changes?

No API change. An In predicate with more than 200 values now prunes files by bounds. Scans, deletes and upserts with large key sets read fewer files.

An IN predicate with more than 200 values returned ROWS_MIGHT_MATCH
before the file bounds were checked. Above that limit, plan_files kept
every data file, so a delete or an upsert scanned the whole table.

Keep the exact check at 200 values or fewer. Above the limit, compare
only the smallest and the largest value against the file bounds. The
result is always a superset of the exact check.

Closes apache#4003
@azwanzuharimi
azwanzuharimi force-pushed the fix-in-predicate-pruning branch from e0b3732 to 930da88 Compare September 30, 2026 17:48
Comment thread pyiceberg/expressions/visitors.py Outdated
Comment on lines +1404 to +1410
if above_limit:
if max(literals) < lower_bound: # type: ignore[operator]
return ROWS_CANNOT_MATCH
else:
literals = {lit for lit in literals if lower_bound <= lit} # type: ignore[operator]
if len(literals) == 0:
return ROWS_CANNOT_MATCH

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.

I need to do a bit more thinking, since the two branches are actually identical:

max({1,2,3}) < 5
vs
len(list(lit for lit in [1,2,3] if 5 <= lit)) == 0

@azwanzuharimi azwanzuharimi Oct 2, 2026 •

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

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

Agreed. The first check gives the same answer in both branches, and both scan every value once. The branches differ only at the upper bound. The filter path narrows the set first, so it prunes a file, when all values sit outside for both sides. The min and max path keeps that file.

So the above limit branch adds no speed and loses pruning, now removed. Same for the manifest visitor: its all(...) checks stop at the first match, so a min and max swap is never faster. Removed too.

The change is now only the removal of the two early returns above the limit. Both visitors run the same check for all sizes. IN_PREDICATE_LIMIT stays in the module because the tests use it, despite no code reads it now.

Measured on a local table with 20 data files of 100,000 rows each. The keys all come from one file. Planning time is near equal. The scan reads fewer files :)

keys files planned before files planned after scan time before scan time after
201 20 / 20 1 / 20 22 ms 12 ms
1,000 20 / 20 1 / 20 39 ms 14 ms
10,000 20 / 20 1 / 20 239 ms 37 ms
100,000 20 / 20 1 / 20 2,234 ms 279 ms
keys plan before plan after
201 8.4 ms 8.9 ms
10,000 23 ms 25 ms
100,000 154 ms 171 ms

The min and max branch gives the same answer as the exact check at the
first bound. Both scan every value once. The exact check also prunes a
file when all values sit outside both bounds. Run the exact check for
all sizes in both visitors.
@azwanzuharimi azwanzuharimi changed the title Prune files on large IN predicates with a min and max check Remove the IN predicate limit so large IN predicates prune files Oct 2, 2026

This branch has not been deployed

No deployments
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.

IN predicates > 200 values disable file pruning

2 participants