-
Notifications
You must be signed in to change notification settings - Fork 4k
[improvement](lance) Push down array membership and boolean scalar index predicates #68687
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: branch-4.1
Are you sure you want to change the base?
Changes from all commits
f46234d
531fcbd
f5b9247
4369546
dac58f5
0b8517a
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,106 @@ | ||
| # Licensed to the Apache Software Foundation (ASF) under one | ||
| # or more contributor license agreements. See the NOTICE file | ||
| # distributed with this work for additional information | ||
| # regarding copyright ownership. The ASF licenses this file | ||
| # to you under the Apache License, Version 2.0 (the | ||
| # "License"); you may not use this file except in compliance | ||
| # with the License. You may obtain a copy of the License at | ||
| # | ||
| # http://www.apache.org/licenses/LICENSE-2.0 | ||
| # | ||
| # Unless required by applicable law or agreed to in writing, | ||
| # software distributed under the License is distributed on an | ||
| # "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY | ||
| # KIND, either express or implied. See the License for the | ||
| # specific language governing permissions and limitations | ||
| # under the License. | ||
|
|
||
| """Build isolated array pushdown fixtures with lance_fixture_requirements.txt. | ||
|
|
||
| The standard MinIO mirror publishes these below warehouse/lance/predicate_arrays. | ||
| Each dataset uses relative data/index paths and two eight-row fragments. The | ||
| partial dataset appends one unindexed fragment after creating its LabelList index. | ||
| """ | ||
|
|
||
| import argparse | ||
| from pathlib import Path | ||
|
|
||
| import lance | ||
| import pyarrow as pa | ||
|
|
||
|
|
||
| LABELS = [["red"], ["blue"], ["red", "blue"], [], None, | ||
| ["red", "red"], [None, "blue"], ["blue", "red"]] * 2 | ||
|
|
||
|
|
||
| def build(output): | ||
| if not __debug__: | ||
| raise RuntimeError("Fixture verification requires assertions; do not use python -O") | ||
| if lance.__version__ != "7.0.0": | ||
| raise RuntimeError("Use the pinned pylance 7.0.0 fixture writer") | ||
| output.mkdir(parents=True, exist_ok=True) | ||
| table = pa.table({"id": pa.array(range(16), type=pa.int64()), | ||
| "labels": pa.array(LABELS, type=pa.list_(pa.string())), | ||
| "category": pa.array([i % 3 for i in range(16)], type=pa.int32())}) | ||
| for name in ["indexed", "partial", "unindexed"]: | ||
| path = output / (name + ".lance") | ||
| if path.exists(): | ||
| raise FileExistsError(path) | ||
| first = table.slice(0, 8) if name == "partial" else table | ||
| dataset = lance.write_dataset(first, str(path), max_rows_per_file=8, max_rows_per_group=8) | ||
| if name != "unindexed": | ||
| dataset.create_scalar_index("labels", "LABEL_LIST", name="labels_idx") | ||
| dataset.create_scalar_index("category", "BTREE", name="category_idx") | ||
| if name == "partial": | ||
| lance.write_dataset(table.slice(8), str(path), mode="append", max_rows_per_file=8) | ||
| check(output) | ||
|
|
||
|
|
||
| def check(output): | ||
| if not __debug__: | ||
| raise RuntimeError("Fixture verification requires assertions; do not use python -O") | ||
| table = pa.table({"id": pa.array(range(16), type=pa.int64()), | ||
| "labels": pa.array(LABELS, type=pa.list_(pa.string())), | ||
| "category": pa.array([i % 3 for i in range(16)], type=pa.int32())}) | ||
| for name in ["indexed", "partial", "unindexed"]: | ||
| path = output / (name + ".lance") | ||
| dataset = lance.dataset(str(path)) | ||
| assert dataset.to_table().to_pydict() == table.to_pydict() | ||
| fragments = dataset.get_fragments() | ||
| assert len(fragments) == 2 | ||
| assert [f.to_table(columns=["id"])["id"].to_pylist() for f in fragments] == [ | ||
| list(range(8)), list(range(8, 16))] | ||
| indexes = {index["name"]: index for index in dataset.list_indices()} | ||
| assert set(indexes) == (set() if name == "unindexed" else {"labels_idx", "category_idx"}) | ||
| expected_coverage = {f.fragment_id for f in (fragments[:1] if name == "partial" else fragments)} | ||
| # Names alone do not prove the mixed indexed/unindexed fixture: verify each | ||
| # physical domain so rebuilding both fragments into an index cannot pass silently. | ||
| for index_name, field, index_type in [("labels_idx", "labels", "LabelList"), | ||
| ("category_idx", "category", "BTree")]: | ||
| if name == "unindexed": | ||
| continue | ||
| index = indexes[index_name] | ||
| assert index["fields"] == [field] and index["type"] == index_type | ||
| assert index["fragment_ids"] == expected_coverage, (name, index_name, index) | ||
| if name == "partial": | ||
| assert len({f.fragment_id for f in fragments} - expected_coverage) == 1 | ||
| for predicate, expected in [ | ||
| ("array_contains(labels, 'red')", [0, 2, 5, 7, 8, 10, 13, 15]), | ||
| ("array_contains(labels, 'red') AND array_contains(labels, 'blue')", [2, 7, 10, 15]), | ||
| ("array_contains(labels, 'red') OR array_contains(labels, 'blue')", | ||
| [0, 1, 2, 5, 6, 7, 8, 9, 10, 13, 14, 15])]: | ||
| for use_index in [False, True]: | ||
| ids = dataset.to_table(columns=["id"], filter=predicate, | ||
| use_scalar_index=use_index)["id"].to_pylist() | ||
| assert sorted(ids) == expected, (name, predicate, ids) | ||
| assert len(dataset.list_indices()) == (0 if name == "unindexed" else 2) | ||
| print("Verified array predicate fixtures (indexed, partial, unindexed)") | ||
|
|
||
|
|
||
| if __name__ == "__main__": | ||
| parser = argparse.ArgumentParser(description=__doc__) | ||
| parser.add_argument("--output", type=Path, default=Path(__file__).parent / | ||
| "preinstalled_data/lance/predicate_arrays") | ||
| parser.add_argument("--check", action="store_true", help="Verify committed fixtures without rebuilding") | ||
| args = parser.parse_args() | ||
| (check if args.check else build)(args.output) | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -89,6 +89,7 @@ | |
| import lance_namespace | ||
| import pyarrow as pa | ||
| import pyarrow.ipc as ipc | ||
| from lance_build_array_predicates import build as build_array_predicates, check as check_array_predicates | ||
| from lance_build_multivector import build as build_multivector, check as check_multivector | ||
| from lance_build_nested_null import build as build_nested_null, check as check_nested_null | ||
| from lance_build_time_travel import check as check_time_travel | ||
|
|
@@ -1671,6 +1672,8 @@ def check_fts_dataset(location: str, *, table_name: str, index_name: str, | |
|
|
||
|
|
||
| def check_catalog(root: Path) -> None: | ||
| # The SQL pushdown suite also relies on these independently published datasets. | ||
| check_array_predicates(root / "predicate_arrays") | ||
| check_data_shapes() | ||
| namespace = lance_namespace.connect("dir", {"root": str(root)}) | ||
| tables = namespace.list_tables(ListTablesRequest(id=[NAMESPACE])) | ||
|
|
@@ -1810,6 +1813,7 @@ def main() -> int: | |
| staging = Path(staging_name) / "lance" | ||
| staging.mkdir() | ||
| build(staging, all_types_source, time_travel_source) | ||
| build_array_predicates(staging / "predicate_arrays") | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [P2] Include predicate_arrays in the committed-fixture self-check. The
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Fixed in 531fcbd. Extracted a reusable array-fixture checker and call it from check_catalog(), so the existing --check path validates all three datasets without rebuilding. The full committed-catalog check passes, and a negative probe confirms missing array datasets are rejected. |
||
| check_catalog(staging) | ||
| backup = output.with_name(output.name + ".old") | ||
| if backup.exists(): | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1 @@ | ||
| {"version":3} |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1 @@ | ||
| {"version":4} |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1 @@ | ||
| {"version":1} |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -17,6 +17,7 @@ | |
|
|
||
| package org.apache.doris.datasource.lance.source; | ||
|
|
||
| import org.apache.doris.analysis.ArrayLiteral; | ||
| import org.apache.doris.analysis.BinaryPredicate; | ||
| import org.apache.doris.analysis.BoolLiteral; | ||
| import org.apache.doris.analysis.CompoundPredicate; | ||
|
|
@@ -48,6 +49,7 @@ | |
| import org.apache.arrow.vector.types.pojo.Field; | ||
| import org.apache.arrow.vector.types.pojo.Schema; | ||
|
|
||
| import java.io.InputStream; | ||
| import java.math.BigDecimal; | ||
| import java.math.BigInteger; | ||
| import java.time.DateTimeException; | ||
|
|
@@ -67,6 +69,8 @@ | |
| /** Converts predicates with identical Doris and Lance semantics to a Substrait ExtendedExpression. */ | ||
| public class LancePredicateConverter { | ||
| private static final TypeCreator REQUIRED = TypeCreator.of(false); | ||
| private static final String LANCE_FUNCTIONS = "https://github.com/apache/doris/blob/master/" | ||
| + "fe/fe-core/src/main/resources/substrait/lance_functions.yaml"; | ||
| private static final SimpleExtension.ExtensionCollection EXTENSIONS = loadExtensions(); | ||
|
|
||
| private final Schema schema; | ||
|
|
@@ -125,7 +129,7 @@ private Optional<Expression> convert(Expr expr) { | |
| return convertLike((LikePredicate) expr); | ||
| } | ||
| if (expr instanceof FunctionCallExpr) { | ||
| return convertStringFunction((FunctionCallExpr) expr); | ||
| return convertFunction((FunctionCallExpr) expr); | ||
| } | ||
| if (expr instanceof SlotRef) { | ||
| return convertBooleanSlot((SlotRef) expr); | ||
|
|
@@ -264,14 +268,18 @@ private Optional<Expression> convertLike(LikePredicate predicate) { | |
| return convertStringPredicate("like:str_str", predicate.getChild(0), predicate.getChild(1), true); | ||
| } | ||
|
|
||
| private Optional<Expression> convertStringFunction(FunctionCallExpr function) { | ||
| private Optional<Expression> convertFunction(FunctionCallExpr function) { | ||
| if (function.getFnName() == null || function.getFn() == null | ||
| || function.getFn().getBinaryType() != TFunctionBinaryType.BUILTIN | ||
| || function.getChildren().size() != 2) { | ||
| return Optional.empty(); | ||
| } | ||
| String functionName = function.getFnName().getFunction().toLowerCase(Locale.ROOT); | ||
| switch (functionName) { | ||
| case "arrays_overlap": | ||
| return convertArraysOverlap(function); | ||
| case "array_contains": | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [P2] Preserve pushdown for three or more OR'ed label checks. Nereids' normal
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Fixed in 531fcbd. Nonempty constant-string arrays_overlap now lowers to a balanced OR of array_has calls, including the three-label Nereids rewrite and either operand order. NULL needles and unsupported forms remain residual. Added FE tests and three-label SQL/Profile cases; the new FE test failed before the fix. Actual FE Substrait integration with upstream main plus Foyer passed, including nullable lists and NOT overlap. Full SQL execution on this revision is pending CI. |
||
| return convertArrayContains(function); | ||
| case "like": | ||
| return convertStringPredicate( | ||
| "like:str_str", function.getChild(0), function.getChild(1), true); | ||
|
|
@@ -286,6 +294,64 @@ private Optional<Expression> convertStringFunction(FunctionCallExpr function) { | |
| } | ||
| } | ||
|
|
||
| private Optional<Expression> convertArraysOverlap(FunctionCallExpr function) { | ||
| Expr input = function.getChild(0); | ||
| Expr values = function.getChild(1); | ||
| if (input instanceof ArrayLiteral) { | ||
| input = function.getChild(1); | ||
| values = function.getChild(0); | ||
| } | ||
| SlotRef slot = directSlot(input); | ||
| ResolvedField field = slot == null ? null : findField(slot); | ||
| if (field == null || !LanceSubstraitSerializer.supportsStringList(field.field) | ||
| || !(values instanceof ArrayLiteral) || values.getChildren().isEmpty()) { | ||
| return Optional.empty(); | ||
| } | ||
| // Nereids folds three or more membership disjuncts into arrays_overlap. Expand | ||
| // only non-NULL string needles: Doris also matches NULL elements to NULL needles. | ||
| List<Expression> level = new ArrayList<>(); | ||
| for (Expr needle : values.getChildren()) { | ||
| if (!(needle instanceof StringLiteral)) { | ||
| return Optional.empty(); | ||
| } | ||
| level.add(arrayMembership(field, needle.getStringValue())); | ||
| } | ||
| // Keep expression depth logarithmic for large constant label lists. | ||
| while (level.size() > 1) { | ||
| List<Expression> next = new ArrayList<>(); | ||
| for (int i = 0; i < level.size(); i += 2) { | ||
| next.add(i + 1 == level.size() ? level.get(i) | ||
| : booleanFunction("or:bool", Arrays.asList(level.get(i), level.get(i + 1)))); | ||
| } | ||
| level = next; | ||
| } | ||
| return Optional.of(level.get(0)); | ||
| } | ||
|
|
||
| private Optional<Expression> convertArrayContains(FunctionCallExpr function) { | ||
| SlotRef slot = directSlot(function.getChild(0)); | ||
| ResolvedField field = slot == null ? null : findField(slot); | ||
| Expr needle = function.getChild(1); | ||
| if (field == null || !LanceSubstraitSerializer.supportsStringList(field.field) | ||
| || !(needle instanceof StringLiteral)) { | ||
| return Optional.empty(); | ||
| } | ||
| // Doris matches a NULL needle against NULL elements; DataFusion does not. Also, | ||
| // array_contains_all is a contiguous subsequence test, not DataFusion's set containment. | ||
| // Only non-NULL scalar membership is interchangeable across the two engines. | ||
| return Optional.of(arrayMembership(field, needle.getStringValue())); | ||
| } | ||
|
|
||
| private Expression arrayMembership(ResolvedField field, String needle) { | ||
| return Expression.ScalarFunctionInvocation.builder() | ||
| .declaration(EXTENSIONS.getScalarFunction( | ||
| SimpleExtension.FunctionAnchor.of(LANCE_FUNCTIONS, "array_has:list_str"))) | ||
| .outputType(TypeCreator.of(field.field.isNullable()).BOOLEAN) | ||
| .arguments(Arrays.asList(fieldReference(field), | ||
| ExpressionCreator.string(false, needle))) | ||
| .build(); | ||
| } | ||
|
|
||
| private Optional<Expression> convertStringPredicate( | ||
| String function, Expr input, Expr pattern, boolean rejectEscapedPattern) { | ||
| SlotRef slot = directSlot(input); | ||
|
|
@@ -497,8 +563,9 @@ private static LiteralExpr directLiteral(Expr expr) { | |
| } | ||
|
|
||
| private static SimpleExtension.ExtensionCollection loadExtensions() { | ||
| try { | ||
| return SimpleExtension.loadDefaults(); | ||
| try (InputStream definitions = | ||
| LancePredicateConverter.class.getResourceAsStream("/substrait/lance_functions.yaml")) { | ||
| return SimpleExtension.loadDefaults().merge(SimpleExtension.load(LANCE_FUNCTIONS, definitions)); | ||
| } catch (Exception e) { | ||
| throw new IllegalStateException("Failed to load Substrait extension definitions", e); | ||
| } | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
[P2] Assert the partial index's actual fragment coverage. The checks here establish two fragments and two index names, but never inspect their
fragment_ids; the partial SQL cases assert a SEGMENT plan and correct rows without checking runtime segment search/fallback counters. If both fragments become indexed, or the segment tasks all fall back, the suite still passes without exercising the intended indexed-plus-unindexed path. Please verify one covered and one uncovered fragment for both indexes and add a partial-table runtime counter assertion.There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Fixed in 531fcbd. Fixture checks now verify the exact fragment_ids for both LabelList and BTree indexes, including exactly one covered and one uncovered fragment in partial. A negative probe rejects an accidentally fully indexed partial fixture. The SQL suite now checks actual segment searches, candidate rows, and zero fallbacks on partial; actual C API integration confirmed those candidate counts.