github-actions[bot] commented on code in PR #68687: URL: https://github.com/apache/doris/pull/68687#discussion_r4147160302
########## docker/thirdparties/docker-compose/iceberg/scripts/lance_build_array_predicates.py: ########## @@ -0,0 +1,76 @@ +# 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) + dataset = lance.dataset(str(path)) + assert dataset.to_table().to_pydict() == table.to_pydict() + assert len(dataset.get_fragments()) == 2 + 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) Review Comment: [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. ########## docker/thirdparties/docker-compose/iceberg/scripts/lance_build_preinstalled_catalog.py: ########## @@ -1810,6 +1811,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") Review Comment: [P2] Include predicate_arrays in the committed-fixture self-check. The `--check` path calls only `check_catalog(output)`, which never opens the new `predicate_arrays` datasets; `build_array_predicates()` runs only during rebuild. A missing or stale committed indexed/partial/unindexed fixture can therefore pass `--check` and fail later in the SQL suite. Please add a reusable check for these datasets to `check_catalog()`. ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LancePredicateConverter.java: ########## @@ -264,14 +267,16 @@ 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 "array_contains": Review Comment: [P2] Preserve pushdown for three or more OR'ed label checks. Nereids' normal `ArrayContainToArrayOverlap` rewrite replaces three distinct `array_contains(labels, ...)` disjuncts with `arrays_overlap(labels, [...])` before this converter runs. This switch leaves `arrays_overlap` residual, so the new LabelList path and segment pruning never run for that common OR shape; the two-label regression cannot catch it. Please handle the rewritten expression with semantics checked against Doris (including NULLs), or keep these disjuncts available to Lance, and add a three-label SQL/profile case. ########## regression-test/suites/external_table_p0/lance/test_lance_array_predicate_pushdown.groovy: ########## @@ -0,0 +1,111 @@ +// 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. + +import org.apache.doris.regression.action.ProfileAction + +suite("test_lance_array_predicate_pushdown", "p0,external") { + String enabled = context.config.otherConfigs.get("enableIcebergTest") + if (enabled == null || !enabled.equalsIgnoreCase("true")) { + logger.info("Lance array pushdown requires the Iceberg MinIO environment") + return + } + String catalog = "test_lance_array_predicate_pushdown" + String endpoint = "http://${context.config.otherConfigs.get('externalEnvIp')}:" + + context.config.otherConfigs.get('iceberg_minio_port') + sql "DROP CATALOG IF EXISTS ${catalog}" + try { + sql """CREATE CATALOG ${catalog} PROPERTIES ( + "type" = "lance", "lance.catalog.type" = "filesystem", + "warehouse" = "s3://warehouse/lance/predicate_arrays", + "s3.endpoint" = "${endpoint}", "s3.access_key" = "admin", + "s3.secret_key" = "password", "s3.region" = "us-east-1", + "use_path_style" = "true")""" + sql "SET enable_profile = true" + sql "SET profile_level = 2" + def profiles = new ProfileAction(context) + String red = "array_contains(labels, 'red')" + String blue = "array_contains(labels, 'blue')" + def cases = [ + [red, [0, 2, 5, 7, 8, 10, 13, 15], 8], + ["${red} AND ${blue}", [2, 7, 10, 15], 4], + ["${red} OR ${blue}", [0, 1, 2, 5, 6, 7, 8, 9, 10, 13, 14, 15], 12], + ["NOT (${red})", [1, 3, 6, 9, 11, 14], null], + ["category IN (0, 1)", [0, 1, 3, 4, 6, 7, 9, 10, 12, 13, 15], 11], + ["category = 0 OR category = 1", [0, 1, 3, 4, 6, 7, 9, 10, 12, 13, 15], 11], + ["${red} AND category = 1", [7, 10, 13], null] + ] + for (String table : ["indexed", "partial", "unindexed"]) { + String relation = "${catalog}.`default`.`${table}`" + cases.eachWithIndex { c, caseId -> + String token = "lance_array_${table}_${caseId}_" + UUID.randomUUID().toString() + String query = "SELECT /* ${token} */ id FROM ${relation} WHERE ${c[0]} ORDER BY id" + explain { + sql(query) + contains "lancePushdownPredicate=" + notContains "predicates:" + if (table != "unindexed") { + contains "lanceScalarIndexScan=SEGMENT" + } else { + notContains "lanceScalarIndexScan=SEGMENT" + } + } + assertEquals(c[1], sql(query).collect { (it[0] as Number).intValue() }) + // A pushed predicate is not necessarily indexed: also verify runtime searches + // and candidate counts, so a non-indexed fallback cannot satisfy this test. + if (table == "indexed" && c[2] != null) { + String profile = profiles.getProfileBySql(token, + ["LanceScalarIndexSegmentsSearched", "LanceScalarIndexCandidateRows"]) + def counter = { String name -> + def matches = profile =~ /${name}: (?:sum )?(\d+)\b/ + assertTrue(matches.find(), "Missing counter ${name}: ${profile}") + return matches.group(1).toLong() + } + assertEquals(1L, counter("LanceScalarIndexSegmentsSearched")) + assertEquals((c[2] as Number).longValue(), counter("LanceScalarIndexCandidateRows")) + assertEquals(0L, counter("LanceScalarIndexSegmentFallbacks")) + } + } + // NULL needle matching and ordered subsequence matching have Doris-specific semantics. + def residualCases = [ + ["array_contains(labels, NULL)", [6, 14]], + ["array_contains_all(labels, ['red', 'blue'])", [2, 10]], + ["array_contains_all(labels, ['red', 'red'])", [5, 13]], + ["${red} OR array_contains(labels, NULL)", [0, 2, 5, 6, 7, 8, 10, 13, 14, 15]] + ] + for (def c : residualCases) { + String query = "SELECT id FROM ${relation} WHERE ${c[0]} ORDER BY id" + explain { + sql(query) + contains "predicates:" + notContains "lancePushdownPredicate=" + } + assertEquals(c[1], sql(query).collect { (it[0] as Number).intValue() }) + } + String mixed = "SELECT id FROM ${relation} WHERE ${blue} AND array_contains(labels, NULL) ORDER BY id" + explain { + sql(mixed) + contains "lancePushdownPredicate=" + contains "predicates:" + } + assertEquals([6, 14], sql(mixed).collect { (it[0] as Number).intValue() }) + // The appended fragment must still contribute matching rows, including with LIMIT. + assertEquals([[15L]], sql("SELECT id FROM ${relation} WHERE ${red} ORDER BY id DESC LIMIT 1")) Review Comment: [P2] Add a query that actually exercises the Lance scanner LIMIT. `ORDER BY id DESC LIMIT 1` becomes TopN; its bound stays on `SortNode`, and the TopN translator forwards a sort limit only to `OlapScanNode`. The Lance scan therefore keeps limit -1 and never sets `TLanceFileDesc.limit`, so this result check cannot detect a regression in indexed or uncovered split-limit handling. Add a fully pushed unordered LIMIT case, assert `lanceLimit=1` in EXPLAIN, and check a selective result on each fragment path. ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScalarIndexPlanner.java: ########## @@ -81,25 +81,38 @@ static Plan plan(LanceTableMetadata metadata, List<Expr> pushedConjuncts, } private static Set<Integer> collectFilterFields(LanceTableMetadata metadata, List<Expr> pushedConjuncts) { - Set<SlotRef> slots = new HashSet<>(); - pushedConjuncts.forEach(expr -> collectDriverSlots(expr, slots)); Set<Integer> fields = new HashSet<>(); - for (SlotRef slot : slots) { - metadata.getLanceFieldId(slot.getColumnName()).ifPresent(fields::add); + for (Expr expr : pushedConjuncts) { + fields.addAll(collectDriverFields(metadata, expr, false)); } return fields; } - private static void collectDriverSlots(Expr expr, Set<SlotRef> slots) { - // A predicate below OR or NOT is not a necessary condition of the whole filter. - // Do not select its index and then force every task into a non-indexed fallback. + private static Set<Integer> collectDriverFields(LanceTableMetadata metadata, Expr expr, boolean complete) { if (expr instanceof CompoundPredicate) { - if (((CompoundPredicate) expr).getOp() == CompoundPredicate.Operator.AND) { - expr.getChildren().forEach(child -> collectDriverSlots(child, slots)); + CompoundPredicate.Operator op = ((CompoundPredicate) expr).getOp(); + if (op == CompoundPredicate.Operator.NOT) { Review Comment: [P2] Preserve scan parallelism for broad NOT results. With one LabelList segment covering many fragments, `NOT array_contains(labels, 'rare')` now selects that segment and `groupFragments` makes one split for all covered fragments; previously NOT produced a split per fragment. If the label is rare or absent, the exact complement still reads nearly every row, but one BE does the work while the other BEs cannot help (the changed scan-node test records four splits becoming one). Please use a cost/selectivity gate to emit fragment splits for unselective complements, or otherwise subdivide the segment domain safely, and cover this plan shape in a test. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
