Remove the IN predicate limit so large IN predicates prune files - #4032
azwanzuharimi wants to merge 2 commits into
Conversation
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
e0b3732 to
930da88
Compare
| 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 |
There was a problem hiding this comment.
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)) == 0There was a problem hiding this comment.
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.
Closes #4003
Rationale for this change
_InclusiveMetricsEvaluationVisitor.visit_inand_ManifestEvalVisitor.visit_inreturnROWS_MIGHT_MATCHas soon as anInpredicate has more thanIN_PREDICATE_LIMIT(200) values. They return before they read any bounds. Soplan_fileskeeps 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_LIMITstays 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.
table/update/validate.pyuses the same evaluator for conflict detection. It drops a file only onROWS_CANNOT_MATCH, so this change can only remove false conflicts.This diverges from Java, which still returns
ROWS_MIGHT_MATCHabove the limit.Are these changes tested?
Yes.
tests/expressions/test_evaluator.py:test_integer_in_above_limit,test_integer_in_at_limitandtest_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_limitfor the manifest evaluator.The above limit tests fail on
main.make lintpasses.tests/expressionsandtests/table/test_evaluator_planning.pypass with 588 tests.Are there any user-facing changes?
No API change. An
Inpredicate with more than 200 values now prunes files by bounds. Scans, deletes and upserts with large key sets read fewer files.