Gabriel39 commented on code in PR #68687:
URL: https://github.com/apache/doris/pull/68687#discussion_r4150148151
##########
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:
Fixed in 531fcbdc462. 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.
##########
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:
Fixed in 531fcbdc462. 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.
##########
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:
Fixed in 531fcbdc462. 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.
##########
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:
Fixed in 531fcbdc462. Complement predicates retain one segment-scoped task
per covered fragment instead of grouping the whole segment into one task. Tasks
share the immutable segment UUID but have disjoint fragment ownership,
preserved on fallback. The FE test now requires four distinct fragment splits
for NOT (it failed before the fix); SQL coverage includes an absent-label
complement and asserts the runtime search count. Disjoint-domain NOT and
NOT-overlap C API integration passed.
##########
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:
Fixed in 531fcbdc462. Added fully pushed, unordered LIMIT 1 queries with
selective predicates targeting the indexed fragment and the uncovered fragment
separately, and assert lanceLimit=1 in EXPLAIN. The existing ordered check
remains as TopN coverage. Actual FE Substrait/C API integration with scanner
limit 1 returns the expected row on both paths; SQL execution is pending CI.
--
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]