diff --git a/.github/workflows/build-thirdparty.yml b/.github/workflows/build-thirdparty.yml index 3d3d11ff27d98a..df6339c8ab81c8 100644 --- a/.github/workflows/build-thirdparty.yml +++ b/.github/workflows/build-thirdparty.yml @@ -80,6 +80,9 @@ jobs: run: | thirdparty/test/adbc-jni-config-test.sh + - name: Test Lance installation + run: bash thirdparty/test/lance-install-test.sh + build_linux: name: Build Third Party Libraries (Linux) needs: changes diff --git a/build.sh b/build.sh index 75179dfac57393..a853802eb629c1 100755 --- a/build.sh +++ b/build.sh @@ -440,6 +440,7 @@ if [[ "${CLEAN}" -eq 1 && "${BUILD_BE}" -eq 0 && "${BUILD_FE}" -eq 0 && ${BUILD_ fi # build thirdparty libraries if necessary. check last thirdparty lib installation +source "${DORIS_HOME}/thirdparty/lance-install.sh" if [[ "${TARGET_SYSTEM}" == 'Darwin' ]]; then LAST_THIRDPARTY_LIB='libbrotlienc.a' else @@ -455,10 +456,21 @@ if [[ ! -f "${DORIS_THIRDPARTY}/installed/lib/${LAST_THIRDPARTY_LIB}" || ! -f "${DORIS_THIRDPARTY}/installed/include/paimon_rust/paimon.h" || ! -s "${DORIS_THIRDPARTY}/installed/lib64/libpaimon_c.a" || ! -s "${DORIS_THIRDPARTY}/installed/include/paimon_rust/paimon.h" || - -e "${DORIS_THIRDPARTY}/installed/lib64/.paimon-installing" ]]; then - # Compilation images may contain only installed artifacts; never erase them without a rebuild source. - if [[ ! -f "${DORIS_THIRDPARTY}/build-thirdparty.sh" ]]; then - echo "Third-party dependencies require a rebuild, but build-thirdparty.sh is missing." >&2 + -e "${DORIS_THIRDPARTY}/installed/lib64/.paimon-installing" ]] || + ! lance_c_install_is_current "${DORIS_HOME}/thirdparty" "${DORIS_THIRDPARTY}/installed"; then + # External trees can be partially updated or pinned to another revision. Preserve + # the existing prefix unless their build inputs can produce the requested Lance version. + for input in build-thirdparty.sh download-thirdparty.sh vars.sh lance-install.sh patches/lance-c-foyer.patch; do + if [[ ! -f "${DORIS_THIRDPARTY}/${input}" || ! -r "${DORIS_THIRDPARTY}/${input}" ]]; then + echo "Third-party dependencies require a rebuild, but ${input} is missing or unreadable." >&2 + echo "Refresh the compilation image or set DORIS_THIRDPARTY to a complete third-party source tree." >&2 + exit 1 + fi + done + if ! expected_lance_fingerprint="$(lance_c_install_fingerprint "${DORIS_HOME}/thirdparty")" || + ! rebuild_lance_fingerprint="$(lance_c_install_fingerprint "${DORIS_THIRDPARTY}")" || + [[ "${rebuild_lance_fingerprint}" != "${expected_lance_fingerprint}" ]]; then + echo "Lance rebuild sources do not match this checkout; installed dependencies have been preserved." >&2 echo "Refresh the compilation image or set DORIS_THIRDPARTY to a complete third-party source tree." >&2 exit 1 fi @@ -471,6 +483,11 @@ if [[ ! -f "${DORIS_THIRDPARTY}/installed/lib/${LAST_THIRDPARTY_LIB}" || else bash "${DORIS_THIRDPARTY}/build-thirdparty.sh" -j "${PARALLEL}" --clean fi + # An external build script can itself be stale. Never link its old output silently. + if ! lance_c_install_is_current "${DORIS_HOME}/thirdparty" "${DORIS_THIRDPARTY}/installed"; then + echo "Lance dependency revision does not match this checkout. Refresh the third-party build tree." >&2 + exit 1 + fi fi update_submodule() { diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/lance_build_array_predicates.py b/docker/thirdparties/docker-compose/iceberg/scripts/lance_build_array_predicates.py new file mode 100644 index 00000000000000..114e5a8a128222 --- /dev/null +++ b/docker/thirdparties/docker-compose/iceberg/scripts/lance_build_array_predicates.py @@ -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) diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/lance_build_preinstalled_catalog.py b/docker/thirdparties/docker-compose/iceberg/scripts/lance_build_preinstalled_catalog.py index 54560857370f72..2645e8748d8337 100644 --- a/docker/thirdparties/docker-compose/iceberg/scripts/lance_build_preinstalled_catalog.py +++ b/docker/thirdparties/docker-compose/iceberg/scripts/lance_build_preinstalled_catalog.py @@ -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") check_catalog(staging) backup = output.with_name(output.name + ".old") if backup.exists(): diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_indices/026150a9-7d9e-4859-a412-fa0582e40071/bitmap_page_lookup.lance b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_indices/026150a9-7d9e-4859-a412-fa0582e40071/bitmap_page_lookup.lance new file mode 100644 index 00000000000000..8b817459d44b7d Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_indices/026150a9-7d9e-4859-a412-fa0582e40071/bitmap_page_lookup.lance differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_indices/79b28d78-5490-40ba-a969-399822c12126/page_data.lance b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_indices/79b28d78-5490-40ba-a969-399822c12126/page_data.lance new file mode 100644 index 00000000000000..af115f4fb6e7b6 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_indices/79b28d78-5490-40ba-a969-399822c12126/page_data.lance differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_indices/79b28d78-5490-40ba-a969-399822c12126/page_lookup.lance b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_indices/79b28d78-5490-40ba-a969-399822c12126/page_lookup.lance new file mode 100644 index 00000000000000..bb72eed316c1cb Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_indices/79b28d78-5490-40ba-a969-399822c12126/page_lookup.lance differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_transactions/0-a4640b8d-55a7-47f8-9392-5ec3d0ed19c9.txn b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_transactions/0-a4640b8d-55a7-47f8-9392-5ec3d0ed19c9.txn new file mode 100644 index 00000000000000..2db5ba7fb9b96b Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_transactions/0-a4640b8d-55a7-47f8-9392-5ec3d0ed19c9.txn differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_transactions/1-d21db0fc-7e46-4699-9948-ab6f8d66098c.txn b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_transactions/1-d21db0fc-7e46-4699-9948-ab6f8d66098c.txn new file mode 100644 index 00000000000000..38fdb4d05345b5 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_transactions/1-d21db0fc-7e46-4699-9948-ab6f8d66098c.txn differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_transactions/2-8498550d-87e2-42c8-931c-1ed0d8bc7371.txn b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_transactions/2-8498550d-87e2-42c8-931c-1ed0d8bc7371.txn new file mode 100644 index 00000000000000..75a68117ca4af1 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_transactions/2-8498550d-87e2-42c8-931c-1ed0d8bc7371.txn differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_versions/18446744073709551612.manifest b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_versions/18446744073709551612.manifest new file mode 100644 index 00000000000000..b96a888dd782c0 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_versions/18446744073709551612.manifest differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_versions/18446744073709551613.manifest b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_versions/18446744073709551613.manifest new file mode 100644 index 00000000000000..4e0e7ac5ed385d Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_versions/18446744073709551613.manifest differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_versions/18446744073709551614.manifest b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_versions/18446744073709551614.manifest new file mode 100644 index 00000000000000..74c137306c73b4 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_versions/18446744073709551614.manifest differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_versions/latest_version_hint.json b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_versions/latest_version_hint.json new file mode 100644 index 00000000000000..f93d3984472d99 --- /dev/null +++ b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/_versions/latest_version_hint.json @@ -0,0 +1 @@ +{"version":3} \ No newline at end of file diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/data/1001010000010111110100106a0464445da0b5a4b7c3465eb5.lance b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/data/1001010000010111110100106a0464445da0b5a4b7c3465eb5.lance new file mode 100644 index 00000000000000..82104af9df6674 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/data/1001010000010111110100106a0464445da0b5a4b7c3465eb5.lance differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/data/111000011101111100111010027ce24037932857b1f2b85e24.lance b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/data/111000011101111100111010027ce24037932857b1f2b85e24.lance new file mode 100644 index 00000000000000..dfe133e40e7628 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/indexed.lance/data/111000011101111100111010027ce24037932857b1f2b85e24.lance differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_indices/9b6b2f0a-bd44-450f-9a3e-ab3f4ace7990/bitmap_page_lookup.lance b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_indices/9b6b2f0a-bd44-450f-9a3e-ab3f4ace7990/bitmap_page_lookup.lance new file mode 100644 index 00000000000000..dfd2305285a29a Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_indices/9b6b2f0a-bd44-450f-9a3e-ab3f4ace7990/bitmap_page_lookup.lance differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_indices/db755621-13aa-4ca9-a495-c0745e34f7c4/page_data.lance b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_indices/db755621-13aa-4ca9-a495-c0745e34f7c4/page_data.lance new file mode 100644 index 00000000000000..7e9051906cd178 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_indices/db755621-13aa-4ca9-a495-c0745e34f7c4/page_data.lance differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_indices/db755621-13aa-4ca9-a495-c0745e34f7c4/page_lookup.lance b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_indices/db755621-13aa-4ca9-a495-c0745e34f7c4/page_lookup.lance new file mode 100644 index 00000000000000..bb72eed316c1cb Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_indices/db755621-13aa-4ca9-a495-c0745e34f7c4/page_lookup.lance differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_transactions/0-58691d36-bff2-4a01-ad25-ae3642019067.txn b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_transactions/0-58691d36-bff2-4a01-ad25-ae3642019067.txn new file mode 100644 index 00000000000000..b8e46a2b19deab Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_transactions/0-58691d36-bff2-4a01-ad25-ae3642019067.txn differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_transactions/1-c343d785-33b7-4448-b0da-160c970f78a1.txn b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_transactions/1-c343d785-33b7-4448-b0da-160c970f78a1.txn new file mode 100644 index 00000000000000..3eb9059d63ee97 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_transactions/1-c343d785-33b7-4448-b0da-160c970f78a1.txn differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_transactions/2-bcdae8c1-73a4-4b25-ad8b-887bf393af0f.txn b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_transactions/2-bcdae8c1-73a4-4b25-ad8b-887bf393af0f.txn new file mode 100644 index 00000000000000..cbf508f445ee5b Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_transactions/2-bcdae8c1-73a4-4b25-ad8b-887bf393af0f.txn differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_transactions/3-8a4b0277-3a05-4eca-8787-2f5ad24763fb.txn b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_transactions/3-8a4b0277-3a05-4eca-8787-2f5ad24763fb.txn new file mode 100644 index 00000000000000..ace146f66c053b Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_transactions/3-8a4b0277-3a05-4eca-8787-2f5ad24763fb.txn differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/18446744073709551611.manifest b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/18446744073709551611.manifest new file mode 100644 index 00000000000000..9e186aac8a1e6f Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/18446744073709551611.manifest differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/18446744073709551612.manifest b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/18446744073709551612.manifest new file mode 100644 index 00000000000000..ca0a927adc0c89 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/18446744073709551612.manifest differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/18446744073709551613.manifest b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/18446744073709551613.manifest new file mode 100644 index 00000000000000..b5cd6bb01d78b7 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/18446744073709551613.manifest differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/18446744073709551614.manifest b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/18446744073709551614.manifest new file mode 100644 index 00000000000000..971179fb5c105c Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/18446744073709551614.manifest differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/latest_version_hint.json b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/latest_version_hint.json new file mode 100644 index 00000000000000..205c7a40a84f18 --- /dev/null +++ b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/_versions/latest_version_hint.json @@ -0,0 +1 @@ +{"version":4} \ No newline at end of file diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/data/011010010011101101000101d214d7401ba969a49532b0ae8b.lance b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/data/011010010011101101000101d214d7401ba969a49532b0ae8b.lance new file mode 100644 index 00000000000000..82104af9df6674 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/data/011010010011101101000101d214d7401ba969a49532b0ae8b.lance differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/data/11001100100100001001011139a9b14af698a84d7cb598d6a7.lance b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/data/11001100100100001001011139a9b14af698a84d7cb598d6a7.lance new file mode 100644 index 00000000000000..dfe133e40e7628 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/partial.lance/data/11001100100100001001011139a9b14af698a84d7cb598d6a7.lance differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/_transactions/0-7c18f76d-3524-4079-835f-7404de102fa9.txn b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/_transactions/0-7c18f76d-3524-4079-835f-7404de102fa9.txn new file mode 100644 index 00000000000000..7bb4aa01baefef Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/_transactions/0-7c18f76d-3524-4079-835f-7404de102fa9.txn differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/_versions/18446744073709551614.manifest b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/_versions/18446744073709551614.manifest new file mode 100644 index 00000000000000..b0e2e53d70671e Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/_versions/18446744073709551614.manifest differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/_versions/latest_version_hint.json b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/_versions/latest_version_hint.json new file mode 100644 index 00000000000000..491d734467aa3c --- /dev/null +++ b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/_versions/latest_version_hint.json @@ -0,0 +1 @@ +{"version":1} \ No newline at end of file diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/data/0011111110010010101001015f6f6345d896c3324132a6d562.lance b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/data/0011111110010010101001015f6f6345d896c3324132a6d562.lance new file mode 100644 index 00000000000000..82104af9df6674 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/data/0011111110010010101001015f6f6345d896c3324132a6d562.lance differ diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/data/1001011111100100010111008b02e34ab199273cee3ff3688e.lance b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/data/1001011111100100010111008b02e34ab199273cee3ff3688e.lance new file mode 100644 index 00000000000000..dfe133e40e7628 Binary files /dev/null and b/docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/predicate_arrays/unindexed.lance/data/1001011111100100010111008b02e34ab199273cee3ff3688e.lance differ diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LancePredicateConverter.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LancePredicateConverter.java index 2e1cd8b890c418..2b23992c9dd19b 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LancePredicateConverter.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LancePredicateConverter.java @@ -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 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,7 +268,7 @@ private Optional convertLike(LikePredicate predicate) { return convertStringPredicate("like:str_str", predicate.getChild(0), predicate.getChild(1), true); } - private Optional convertStringFunction(FunctionCallExpr function) { + private Optional convertFunction(FunctionCallExpr function) { if (function.getFnName() == null || function.getFn() == null || function.getFn().getBinaryType() != TFunctionBinaryType.BUILTIN || function.getChildren().size() != 2) { @@ -272,6 +276,10 @@ private Optional convertStringFunction(FunctionCallExpr function) { } String functionName = function.getFnName().getFunction().toLowerCase(Locale.ROOT); switch (functionName) { + case "arrays_overlap": + return convertArraysOverlap(function); + case "array_contains": + return convertArrayContains(function); case "like": return convertStringPredicate( "like:str_str", function.getChild(0), function.getChild(1), true); @@ -286,6 +294,64 @@ private Optional convertStringFunction(FunctionCallExpr function) { } } + private Optional 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 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 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 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 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); } diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScalarIndexPlanner.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScalarIndexPlanner.java index ab259a05fe2489..4d6319242f6984 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScalarIndexPlanner.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScalarIndexPlanner.java @@ -17,9 +17,17 @@ package org.apache.doris.datasource.lance.source; +import org.apache.doris.analysis.ArrayLiteral; +import org.apache.doris.analysis.BinaryPredicate; import org.apache.doris.analysis.CompoundPredicate; import org.apache.doris.analysis.Expr; +import org.apache.doris.analysis.FunctionCallExpr; +import org.apache.doris.analysis.InPredicate; +import org.apache.doris.analysis.IsNullPredicate; +import org.apache.doris.analysis.LikePredicate; +import org.apache.doris.analysis.LiteralExpr; import org.apache.doris.analysis.SlotRef; +import org.apache.doris.analysis.StringLiteral; import org.apache.doris.datasource.lance.index.LanceIndexSegmentGroup; import org.apache.doris.datasource.lance.index.LanceIndexSegmentInfo; import org.apache.doris.datasource.lance.metadata.LanceFragmentInfo; @@ -28,6 +36,7 @@ import org.lance.index.IndexType; import java.util.ArrayList; +import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; @@ -35,6 +44,11 @@ /** Assigns one BTree/Bitmap/LabelList segment and a disjoint fragment domain to each ordinary scan task. */ final class LanceScalarIndexPlanner { + // Match the pinned lance-c scoped-expression bounds. Exceeding either bound makes + // native planning scan the whole domain, so do not coalesce fragments in that case. + private static final int MAX_EXPRESSION_NODES = 128; + private static final int MAX_EXPRESSION_DEPTH = 32; + static final class Plan { final String indexName; final LanceSplitBuilder splits; @@ -53,23 +67,39 @@ static Plan plan(LanceTableMetadata metadata, List pushedConjuncts, || !metadata.getIndexMetadataState().canPlanIndexSegments()) { return null; } - Set filterFields = collectFilterFields(metadata, pushedConjuncts); - if (filterFields.isEmpty()) { + if (exceedsExpressionBudget(pushedConjuncts)) { return null; } // Metadata already groups physical segments by logical index. Name order // provides a stable winner when multiple indices cover the same number of rows. List indices = new ArrayList<>(metadata.getIndexes()); indices.sort(java.util.Comparator.comparing(LanceIndexSegmentGroup::getName)); + // Native chooses the first matching parser, not the index with the widest + // coverage. Without its dispatch order, leave competing logical indexes to native. + Set indexedFields = new HashSet<>(); + Set ambiguousFields = new HashSet<>(); + for (LanceIndexSegmentGroup index : indices) { + for (Integer field : index.getSegments().get(0).getFieldIds()) { + if (!indexedFields.add(field)) { + ambiguousFields.add(field); + } + } + } Plan selected = null; + Map> fieldsByIndexType = new HashMap<>(); for (LanceIndexSegmentGroup logicalIndex : indices) { List segments = logicalIndex.getSegments(); - // PR #79 supports one top-level key in BTree/Bitmap/LabelList indices. Lance - // performs the final typed driver selection and falls back within the same domain. + // Only coalesce a segment when a positive, indexable necessary condition + // exists. Broad complements keep fragment parallelism without repeating a + // segment-wide index search once per fragment. LanceIndexSegmentInfo index = segments.get(0); if ((index.getIndexType() != IndexType.BTREE && index.getIndexType() != IndexType.BITMAP && index.getIndexType() != IndexType.LABEL_LIST) - || index.getFieldIds().size() != 1 || !filterFields.contains(index.getFieldIds().get(0))) { + || index.getFieldIds().size() != 1 + || ambiguousFields.contains(index.getFieldIds().get(0)) + || !fieldsByIndexType.computeIfAbsent(index.getIndexType(), + type -> collectFilterFields(metadata, pushedConjuncts, type)) + .contains(index.getFieldIds().get(0))) { continue; } Plan candidate = groupFragments(metadata, segments, visibleFragments); @@ -80,26 +110,199 @@ static Plan plan(LanceTableMetadata metadata, List pushedConjuncts, return selected; } - private static Set collectFilterFields(LanceTableMetadata metadata, List pushedConjuncts) { - Set slots = new HashSet<>(); - pushedConjuncts.forEach(expr -> collectDriverSlots(expr, slots)); + private static Set collectFilterFields(LanceTableMetadata metadata, + List pushedConjuncts, IndexType indexType) { Set fields = new HashSet<>(); + for (Expr expr : pushedConjuncts) { + fields.addAll(collectDriverFields(metadata, expr, indexType, false)); + } + return fields; + } + + private static Set collectDriverFields(LanceTableMetadata metadata, Expr expr, + IndexType indexType, boolean requireExact) { + Set fields = new HashSet<>(); + if (expr instanceof CompoundPredicate) { + CompoundPredicate.Operator op = ((CompoundPredicate) expr).getOp(); + if (op == CompoundPredicate.Operator.NOT) { + return fields; + } + boolean exactBranches = requireExact || op == CompoundPredicate.Operator.OR; + Set left = collectDriverFields(metadata, expr.getChild(0), indexType, exactBranches); + Set right = collectDriverFields(metadata, expr.getChild(1), indexType, exactBranches); + if (op == CompoundPredicate.Operator.AND) { + // A refine-only conjunct anywhere below OR makes native reject the + // union, even if its sibling could independently drive this index. + if (requireExact && (left.isEmpty() || right.isEmpty())) { + return fields; + } + left.addAll(right); + } else { + // Sharing a slot is insufficient: both OR branches must actually be + // indexable. A suffix LIKE, for example, cannot supply BTree candidates. + left.retainAll(right); + } + return left; + } + if (!isPositiveIndexLeaf(expr, indexType, requireExact)) { + return fields; + } + Set slots = new HashSet<>(); + expr.collect(SlotRef.class, slots); for (SlotRef slot : slots) { metadata.getLanceFieldId(slot.getColumnName()).ifPresent(fields::add); } return fields; } - private static void collectDriverSlots(Expr expr, Set 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 boolean isPositiveIndexLeaf(Expr expr, IndexType indexType, boolean requireExact) { + if (indexType == IndexType.LABEL_LIST) { + if (!(expr instanceof FunctionCallExpr) || expr.getChildren().size() != 2) { + return false; + } + String name = ((FunctionCallExpr) expr).getFnName().getFunction(); + if ("array_contains".equalsIgnoreCase(name)) { + return expr.getChild(0) instanceof SlotRef && expr.getChild(1) instanceof LiteralExpr; + } + return "arrays_overlap".equalsIgnoreCase(name) && overlapSize(expr) > 0; + } + if (expr instanceof BinaryPredicate) { + BinaryPredicate.Operator op = ((BinaryPredicate) expr).getOp(); + return (op == BinaryPredicate.Operator.EQ || op == BinaryPredicate.Operator.GT + || op == BinaryPredicate.Operator.GE || op == BinaryPredicate.Operator.LT + || op == BinaryPredicate.Operator.LE) + && ((expr.getChild(0) instanceof SlotRef && expr.getChild(1) instanceof LiteralExpr) + || (expr.getChild(1) instanceof SlotRef && expr.getChild(0) instanceof LiteralExpr)); + } + if (expr instanceof InPredicate) { + return !((InPredicate) expr).isNotIn() && expr.getChild(0) instanceof SlotRef; + } + if (expr instanceof IsNullPredicate) { + return !((IsNullPredicate) expr).isNotNull() && expr.getChild(0) instanceof SlotRef; + } + // Bitmap has no prefix-query support. A refined LIKE can drive an AND, + // but native's OR planner rejects branches that need a residual recheck. + if (indexType != IndexType.BTREE || expr.getChildren().size() != 2 + || !(expr.getChild(0) instanceof SlotRef) || !(expr.getChild(1) instanceof StringLiteral)) { + return false; + } + String name = expr instanceof FunctionCallExpr + ? ((FunctionCallExpr) expr).getFnName().getFunction() : ""; + String prefix = ((StringLiteral) expr.getChild(1)).getStringValue(); + if ("starts_with".equalsIgnoreCase(name)) { + // Unlike LIKE, every character in starts_with is literal, including % and _. + return !prefix.isEmpty(); + } + boolean like = (expr instanceof LikePredicate && ((LikePredicate) expr).getOp() == LikePredicate.Operator.LIKE) + || "like".equalsIgnoreCase(name); + if (!like || prefix.isEmpty() || prefix.indexOf('\\') >= 0) { + return false; + } + for (int i = 0; i < prefix.length(); i++) { + char character = prefix.charAt(i); + if (character == '%' || character == '_') { + return i > 0 && (!requireExact || (character == '%' && i == prefix.length() - 1)); + } + } + return true; + } + + private static int overlapSize(Expr expr) { + if (!(expr instanceof FunctionCallExpr) || expr.getChildren().size() != 2 + || !"arrays_overlap".equalsIgnoreCase(((FunctionCallExpr) expr).getFnName().getFunction())) { + return 0; + } + for (int i = 0; i < 2; i++) { + if (expr.getChild(i) instanceof ArrayLiteral && expr.getChild(1 - i) instanceof SlotRef) { + return expr.getChild(i).getChildren().size(); + } + } + return 0; + } + + static boolean shouldDisableFragmentIndex(List pushedConjuncts) { + if (pushedConjuncts.isEmpty()) { + return false; + } + // Missing metadata or an ambiguous index name is not evidence against native + // index use. Disable only known expensive shapes, independent of FE discovery. + return exceedsExpressionBudget(pushedConjuncts) + || pushedConjuncts.stream().noneMatch(LanceScalarIndexPlanner::hasPositivePredicate); + } + + private static boolean hasPositivePredicate(Expr expr) { if (expr instanceof CompoundPredicate) { - if (((CompoundPredicate) expr).getOp() == CompoundPredicate.Operator.AND) { - expr.getChildren().forEach(child -> collectDriverSlots(child, slots)); + return ((CompoundPredicate) expr).getOp() != CompoundPredicate.Operator.NOT + && expr.getChildren().stream().anyMatch(LanceScalarIndexPlanner::hasPositivePredicate); + } + if (expr instanceof BinaryPredicate) { + return ((BinaryPredicate) expr).getOp() != BinaryPredicate.Operator.NE; + } + if (expr instanceof InPredicate) { + return !((InPredicate) expr).isNotIn(); + } + if (expr instanceof IsNullPredicate) { + return !((IsNullPredicate) expr).isNotNull(); + } + return true; + } + + private static boolean exceedsExpressionBudget(List conjuncts) { + return conjunctionNodes(conjuncts, 0, conjuncts.size(), 0) > MAX_EXPRESSION_NODES; + } + + private static int conjunctionNodes(List conjuncts, int begin, int end, int depth) { + if (begin == end) { + return 0; + } + if (end - begin == 1) { + return expressionNodes(conjuncts.get(begin), depth); + } + // The converter emits n-ary and:bool; DataFusion splits its arguments in + // halves. Preserve each leaf's actual depth instead of assuming a left-deep AND. + int middle = begin + (end - begin) / 2; + return 1 + conjunctionNodes(conjuncts, begin, middle, depth + 1) + + conjunctionNodes(conjuncts, middle, end, depth + 1); + } + + private static int expressionNodes(Expr expr, int depth) { + if (depth > MAX_EXPRESSION_DEPTH) { + return MAX_EXPRESSION_NODES + 1; + } + if (expr instanceof CompoundPredicate) { + int nodes = 1; + for (Expr child : expr.getChildren()) { + nodes += expressionNodes(child, depth + 1); + if (nodes > MAX_EXPRESSION_NODES) { + break; + } + } + return nodes; + } + int labels = overlapSize(expr); + // The converter emits a balanced OR tree: N memberships require 2*N-1 nodes. + if (labels > 0) { + int levels = 32 - Integer.numberOfLeadingZeros(labels - 1); + if (labels > MAX_EXPRESSION_NODES / 2 || depth + levels > MAX_EXPRESSION_DEPTH) { + return MAX_EXPRESSION_NODES + 1; + } + return 2 * labels - 1; + } + if (expr instanceof InPredicate) { + InPredicate in = (InPredicate) expr; + int values = in.getInElementNum(); + // DataFusion 54 expands up to three values into a left-deep OR (or + // AND of negated equalities for NOT IN) before the native budget check. + if (values > 0 && values <= 3) { + int negation = in.isNotIn() ? 1 : 0; + if (depth + values - 1 + negation > MAX_EXPRESSION_DEPTH) { + return MAX_EXPRESSION_NODES + 1; + } + return (2 + negation) * values - 1; } - } else { - expr.collect(SlotRef.class, slots); + return in.isNotIn() ? 2 : 1; } + return 1; } private static Plan groupFragments(LanceTableMetadata metadata, List segments, diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScanPlanner.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScanPlanner.java index e95420582b6e57..23fc3eeb5ddbcd 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScanPlanner.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScanPlanner.java @@ -195,7 +195,11 @@ private List createNormalFragmentSplits(LanceTableMetadata metadata, } else { plan = new LanceSplitBuilder(metadata.getDatasetUri(), metadata.getVersion(), 0); } - plan.addUncoveredFragments(visibleFragments.values(), 1, scalarIndexPlan != null); + // Uncovered fragments must not repeat the selected segment's search. When FE + // could not select a segment, retain native selection unless the filter itself + // is known to cause expensive repeated searches (a complement or oversized tree). + plan.addUncoveredFragments(visibleFragments.values(), 1, scalarIndexPlan != null + || LanceScalarIndexPlanner.shouldDisableFragmentIndex(lancePushedConjuncts)); return plan.buildSplits(); } diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceSubstraitSerializer.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceSubstraitSerializer.java index e5323bde81535d..79d3f4792617a0 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceSubstraitSerializer.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceSubstraitSerializer.java @@ -108,10 +108,19 @@ static boolean supportsType(ArrowType type) { || type instanceof ArrowType.LargeUtf8; } + static boolean supportsStringList(Field field) { + // Keep this separate from scalar comparability: list serialization must not enable + // array comparisons, IN, or null checks without their own semantic coverage. + return field.getType() instanceof ArrowType.List && field.getChildren().size() == 1 + && field.getChildren().get(0).getType() instanceof ArrowType.Utf8; + } + static Type fieldType(Field field) { TypeCreator creator = TypeCreator.of(field.isNullable()); ArrowType type = field.getType(); - if (type instanceof ArrowType.Bool) { + if (supportsStringList(field)) { + return creator.list(fieldType(field.getChildren().get(0))); + } else if (type instanceof ArrowType.Bool) { return creator.BOOLEAN; } else if (type instanceof ArrowType.Int) { switch (((ArrowType.Int) type).getBitWidth()) { @@ -198,6 +207,13 @@ byte[] serialize(Expression expression) { } private static Optional toSubstraitProtoType(Field field) { + if (supportsStringList(field)) { + return Optional.of(io.substrait.proto.Type.newBuilder() + .setList(io.substrait.proto.Type.List.newBuilder() + .setNullability(nullability(field)) + .setType(toSubstraitProtoType(field.getChildren().get(0)).get())) + .build()); + } if (!supportsType(field.getType())) { return Optional.empty(); } diff --git a/fe/fe-core/src/main/resources/substrait/lance_functions.yaml b/fe/fe-core/src/main/resources/substrait/lance_functions.yaml new file mode 100644 index 00000000000000..3e6f64ecbb99e4 --- /dev/null +++ b/fe/fe-core/src/main/resources/substrait/lance_functions.yaml @@ -0,0 +1,27 @@ +# 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. + +%YAML 1.2 +--- +scalar_functions: + - name: array_has + description: Non-null UTF-8 label membership in a List column. + impls: + - args: + - value: list + - value: string + return: boolean diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/LancePredicateConverterTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/LancePredicateConverterTest.java index 64e6fd70724430..148fbf1e53dbf5 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/LancePredicateConverterTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/LancePredicateConverterTest.java @@ -17,6 +17,7 @@ package org.apache.doris.datasource.lance; +import org.apache.doris.analysis.ArrayLiteral; import org.apache.doris.analysis.BinaryPredicate; import org.apache.doris.analysis.CompoundPredicate; import org.apache.doris.analysis.DateLiteral; @@ -29,8 +30,10 @@ import org.apache.doris.analysis.IsNullPredicate; import org.apache.doris.analysis.LargeIntLiteral; import org.apache.doris.analysis.LikePredicate; +import org.apache.doris.analysis.NullLiteral; import org.apache.doris.analysis.SlotRef; import org.apache.doris.analysis.StringLiteral; +import org.apache.doris.catalog.ArrayType; import org.apache.doris.catalog.ScalarFunction; import org.apache.doris.catalog.ScalarType; import org.apache.doris.catalog.Type; @@ -73,6 +76,118 @@ public class LancePredicateConverterTest { Field.nullable("timestamp_ns", new ArrowType.Timestamp(TimeUnit.NANOSECOND, null)), Field.nullable("timestamp_us_utc", new ArrowType.Timestamp(TimeUnit.MICROSECOND, "UTC"))))); + @Test + public void testStringArrayMembershipSchemaAndFunction() throws Exception { + // An unsupported preceding field must not shift the array's field reference. + Schema schema = new Schema(Arrays.asList( + Field.nullable("nested", ArrowType.Struct.INSTANCE), + new Field("labels", FieldType.nullable(ArrowType.List.INSTANCE), + Collections.singletonList(Field.nullable("item", ArrowType.Utf8.INSTANCE))))); + Expr predicate = arrayFunction("array_contains", new SlotRef(null, "labels"), new StringLiteral("red")); + LancePredicateConverter.ConversionResult result = + new LancePredicateConverter(schema).convert(Collections.singletonList(predicate)); + Assertions.assertEquals(1, result.getPushedConjuncts().size()); + Assertions.assertTrue(result.getResidualConjuncts().isEmpty()); + ExtendedExpression envelope = ExtendedExpression.parseFrom(result.getSubstraitFilter()); + Assertions.assertEquals(Arrays.asList("__unlikely_name_placeholder_doris_0", "labels"), + envelope.getBaseSchema().getNamesList()); + io.substrait.proto.Type.List list = envelope.getBaseSchema().getStruct().getTypes(1).getList(); + Assertions.assertEquals(io.substrait.proto.Type.Nullability.NULLABILITY_NULLABLE, list.getNullability()); + Assertions.assertTrue(list.getType().hasString()); + Assertions.assertEquals(io.substrait.proto.Type.Nullability.NULLABILITY_NULLABLE, + list.getType().getString().getNullability()); + io.substrait.proto.Expression.ScalarFunction function = + envelope.getReferredExpr(0).getExpression().getScalarFunction(); + Assertions.assertEquals(1, function.getArguments(0).getValue().getSelection() + .getDirectReference().getStructField().getField()); + Assertions.assertEquals("red", function.getArguments(1).getValue().getLiteral().getString()); + Assertions.assertEquals(io.substrait.proto.Type.Nullability.NULLABILITY_NULLABLE, + function.getOutputType().getBool().getNullability()); + Assertions.assertTrue(envelope.getExtensionsList().stream() + .anyMatch(extension -> extension.hasExtensionFunction() + && extension.getExtensionFunction().getName().startsWith("array_has:"))); + } + + @Test + public void testArrayMembershipBooleanCombinations() { + LancePredicateConverter arrays = arrayConverter(ArrowType.List.INSTANCE, ArrowType.Utf8.INSTANCE); + Expr red = arrayFunction("array_contains", new SlotRef(null, "labels"), new StringLiteral("red")); + Expr blue = arrayFunction("array_contains", new SlotRef(null, "labels"), new StringLiteral("blue")); + for (Expr predicate : Arrays.asList(red, + new CompoundPredicate(CompoundPredicate.Operator.AND, red, blue), + new CompoundPredicate(CompoundPredicate.Operator.OR, red, blue), + new CompoundPredicate(CompoundPredicate.Operator.NOT, red, null))) { + Assertions.assertEquals(1, arrays.convert(Collections.singletonList(predicate)) + .getPushedConjuncts().size()); + } + } + + @Test + public void testRewrittenArrayOverlapPushdown() throws Exception { + LancePredicateConverter arrays = arrayConverter(ArrowType.List.INSTANCE, ArrowType.Utf8.INSTANCE); + Expr labels = new SlotRef(null, "labels"); + Expr values = new ArrayLiteral(ArrayType.create(Type.STRING, true), + new StringLiteral("red"), new StringLiteral("blue"), new StringLiteral("green")); + for (Expr predicate : Arrays.asList(arrayFunction("arrays_overlap", labels, values), + arrayFunction("arrays_overlap", values, labels), + new CompoundPredicate(CompoundPredicate.Operator.NOT, + arrayFunction("arrays_overlap", labels, values), null))) { + LancePredicateConverter.ConversionResult result = arrays.convert(Collections.singletonList(predicate)); + Assertions.assertTrue(result.getResidualConjuncts().isEmpty()); + String encoded = ExtendedExpression.parseFrom(result.getSubstraitFilter()).toString(); + Assertions.assertTrue(encoded.contains("array_has:list_str")); + Assertions.assertTrue(encoded.contains("or:bool")); + Assertions.assertTrue(encoded.contains("green")); + } + // A NULL needle can match a NULL element in Doris and must not become array_has. + for (Expr valuesWithDifferentSemantics : Arrays.asList( + new ArrayLiteral(ArrayType.create(Type.STRING, true), new StringLiteral("red"), new NullLiteral()), + new ArrayLiteral(), new SlotRef(null, "labels"))) { + Expr predicate = arrayFunction("arrays_overlap", labels, valuesWithDifferentSemantics); + Assertions.assertEquals(Collections.singletonList(predicate), + arrays.convert(Collections.singletonList(predicate)).getResidualConjuncts()); + } + } + + @Test + public void testUnsafeArrayMembershipRemainsResidual() { + LancePredicateConverter arrays = arrayConverter(ArrowType.List.INSTANCE, ArrowType.Utf8.INSTANCE); + Expr red = arrayFunction("array_contains", new SlotRef(null, "labels"), new StringLiteral("red")); + Expr nullNeedle = arrayFunction("array_contains", new SlotRef(null, "labels"), new NullLiteral()); + FunctionCallExpr udf = arrayFunction("array_contains", new SlotRef(null, "labels"), new StringLiteral("red")); + udf.getFn().setBinaryType(TFunctionBinaryType.JAVA_UDF); + for (Expr predicate : Arrays.asList(nullNeedle, udf, + arrayFunction("array_contains_all", new SlotRef(null, "labels"), new StringLiteral("red")), + arrayFunction("array_contains", new SlotRef(null, "labels"), new SlotRef(null, "needle")), + new CompoundPredicate(CompoundPredicate.Operator.OR, red, nullNeedle), + new CompoundPredicate(CompoundPredicate.Operator.NOT, nullNeedle, null), + new IsNullPredicate(new SlotRef(null, "labels"), false))) { + Assertions.assertEquals(Collections.singletonList(predicate), + arrays.convert(Collections.singletonList(predicate)).getResidualConjuncts()); + } + Assertions.assertEquals(1, arrays.convert(Arrays.asList(red, nullNeedle)).getPushedConjuncts().size()); + for (LancePredicateConverter unsupported : Arrays.asList( + arrayConverter(ArrowType.List.INSTANCE, new ArrowType.Int(64, true)), + arrayConverter(ArrowType.List.INSTANCE, ArrowType.LargeUtf8.INSTANCE), + arrayConverter(ArrowType.LargeList.INSTANCE, ArrowType.Utf8.INSTANCE), + arrayConverter(new ArrowType.FixedSizeList(2), ArrowType.Utf8.INSTANCE))) { + Assertions.assertEquals(1, unsupported.convert(Collections.singletonList(red)).getResidualConjuncts().size()); + } + } + + private static LancePredicateConverter arrayConverter(ArrowType container, ArrowType element) { + return new LancePredicateConverter(new Schema(Collections.singletonList( + new Field("labels", FieldType.nullable(container), + Collections.singletonList(Field.nullable("item", element)))))); + } + + private static FunctionCallExpr arrayFunction(String name, Expr... arguments) { + FunctionCallExpr function = new FunctionCallExpr(name, Arrays.asList(arguments)); + function.setFn(new ScalarFunction(new FunctionName(name), + Arrays.asList(ArrayType.create(Type.STRING, true), Type.STRING), Type.BOOLEAN, false, true)); + return function; + } + @Test public void testComparisonAndStringLiteralEncoding() throws Exception { Expr rowId = new SlotRef(null, "row_id"); diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/source/LanceScalarIndexPlannerTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/source/LanceScalarIndexPlannerTest.java new file mode 100644 index 00000000000000..d94e2f554740ef --- /dev/null +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/source/LanceScalarIndexPlannerTest.java @@ -0,0 +1,292 @@ +// 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. + +package org.apache.doris.datasource.lance.source; + +import org.apache.doris.analysis.ArrayLiteral; +import org.apache.doris.analysis.BinaryPredicate; +import org.apache.doris.analysis.CompoundPredicate; +import org.apache.doris.analysis.Expr; +import org.apache.doris.analysis.FunctionCallExpr; +import org.apache.doris.analysis.FunctionName; +import org.apache.doris.analysis.InPredicate; +import org.apache.doris.analysis.IntLiteral; +import org.apache.doris.analysis.LikePredicate; +import org.apache.doris.analysis.SlotRef; +import org.apache.doris.analysis.StringLiteral; +import org.apache.doris.catalog.ArrayType; +import org.apache.doris.catalog.ScalarFunction; +import org.apache.doris.catalog.Type; +import org.apache.doris.datasource.lance.index.LanceIndexSegmentInfo; +import org.apache.doris.datasource.lance.metadata.LanceFragmentInfo; +import org.apache.doris.datasource.lance.metadata.LanceTableAccess; +import org.apache.doris.datasource.lance.metadata.LanceTableMetadata; + +import org.apache.arrow.vector.types.pojo.ArrowType; +import org.apache.arrow.vector.types.pojo.Field; +import org.apache.arrow.vector.types.pojo.FieldType; +import org.apache.arrow.vector.types.pojo.Schema; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.lance.index.IndexType; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.UUID; + +public class LanceScalarIndexPlannerTest { + @Test + public void testBooleanDriverSelectionPreservesAllBranches() { + Expr key = equal("key", 1); + Expr otherKey = equal("key", 2); + Expr residual = equal("other", 3); + for (Expr filter : Arrays.asList(or(key, otherKey), + or(and(key, residual), and(otherKey, residual)), and(not(and(key, residual)), key))) { + LanceScalarIndexPlanner.Plan plan = plan(filter); + Assertions.assertNotNull(plan, filter.toSql()); + Assertions.assertEquals("key_idx", plan.indexName); + Assertions.assertEquals(1, plan.splits.splitCount()); + Assertions.assertTrue(plan.splits.isCoveredByIndexSegment(1)); + Assertions.assertTrue(plan.splits.isCoveredByIndexSegment(2)); + } + // Pruning either an OR branch or an AND below NOT could drop matching rows. + for (Expr filter : Arrays.asList(not(key), not(or(key, otherKey)), or(key, residual), not(and(key, residual)), + not(or(key, residual)), or(and(key, residual), residual))) { + Assertions.assertNull(plan(filter), filter.toSql()); + } + } + + @Test + public void testConvertedArrayPredicatesSelectLabelList() { + Schema schema = new Schema(Collections.singletonList(new Field("labels", + FieldType.nullable(ArrowType.List.INSTANCE), + Collections.singletonList(Field.nullable("item", ArrowType.Utf8.INSTANCE))))); + LanceFragmentInfo fragment = new LanceFragmentInfo(1, 10, 10); + LanceTableMetadata metadata = LanceTableMetadata.createSnapshotWithIndexes( + new LanceTableAccess("s3://bucket/labels.lance", Collections.emptyMap()), 42, schema, + Collections.singletonList(fragment), Collections.singletonMap("labels", 9), + Collections.singletonList(new LanceIndexSegmentInfo(UUID.randomUUID(), "labels_idx", + Collections.singletonList(9), Collections.singletonList(1L), IndexType.LABEL_LIST, null))); + FunctionCallExpr red = contains("red"); + FunctionCallExpr blue = contains("blue"); + for (Expr filter : Arrays.asList(red, and(red, blue), or(red, blue))) { + LancePredicateConverter.ConversionResult converted = + new LancePredicateConverter(schema).convert(Collections.singletonList(filter)); + Assertions.assertTrue(converted.getResidualConjuncts().isEmpty()); + LanceScalarIndexPlanner.Plan plan = LanceScalarIndexPlanner.plan(metadata, + converted.getPushedConjuncts(), Collections.singletonMap(1L, fragment)); + Assertions.assertNotNull(plan); + Assertions.assertEquals("labels_idx", plan.indexName); + Assertions.assertEquals(1, plan.splits.splitCount()); + } + } + + @Test + public void testPositiveDriverWithComplementSearchesSegmentOnce() { + Expr key = equal("key", 1); + Expr otherNotEqual = new BinaryPredicate(BinaryPredicate.Operator.NE, + new SlotRef(null, "other"), new IntLiteral(0)); + Assertions.assertEquals(1, plan(Arrays.asList(key, otherNotEqual), false).splits.splitCount()); + Assertions.assertEquals(1, plan(and(key, otherNotEqual)).splits.splitCount()); + Assertions.assertEquals(1, plan(and(key, not(equal("key", 2)))).splits.splitCount()); + } + + @Test + public void testUnindexableOrBranchDoesNotGroupFragments() { + Expr suffix = new LikePredicate(LikePredicate.Operator.LIKE, + new SlotRef(null, "key"), new StringLiteral("%y%")); + Expr key = new BinaryPredicate(BinaryPredicate.Operator.EQ, + new SlotRef(null, "key"), new StringLiteral("x")); + Assertions.assertNull(plan(Collections.singletonList(or(key, suffix)), true)); + Assertions.assertNotNull(plan(Collections.singletonList(and(key, suffix)), true)); + Expr prefix = new LikePredicate(LikePredicate.Operator.LIKE, + new SlotRef(null, "key"), new StringLiteral("y%")); + Assertions.assertEquals(1, plan(Collections.singletonList(or(key, prefix)), true).splits.splitCount()); + } + + @Test + public void testExpressionDepthBudget() { + Expr filter = equal("key", 1); + for (int depth = 1; depth <= 32; depth++) { + filter = and(filter, equal("key", depth)); + } + Assertions.assertNotNull(plan(filter)); + Assertions.assertNull(plan(and(filter, equal("key", 33)))); + } + + @Test + public void testOverlapExpressionBudget() throws Exception { + Schema schema = new Schema(Collections.singletonList(new Field("labels", + FieldType.nullable(ArrowType.List.INSTANCE), + Collections.singletonList(Field.nullable("item", ArrowType.Utf8.INSTANCE))))); + LanceFragmentInfo first = new LanceFragmentInfo(1, 10, 10); + LanceFragmentInfo second = new LanceFragmentInfo(2, 10, 10); + Map fragments = new HashMap<>(); + fragments.put(1L, first); + fragments.put(2L, second); + LanceTableMetadata metadata = LanceTableMetadata.createSnapshotWithIndexes( + new LanceTableAccess("s3://bucket/labels.lance", Collections.emptyMap()), 42, schema, + Arrays.asList(first, second), Collections.singletonMap("labels", 9), + Collections.singletonList(new LanceIndexSegmentInfo(UUID.randomUUID(), "labels_idx", + Collections.singletonList(9), Arrays.asList(1L, 2L), IndexType.LABEL_LIST, null))); + for (int count : Arrays.asList(64, 65)) { + StringLiteral[] labels = new StringLiteral[count]; + for (int i = 0; i < count; i++) { + labels[i] = new StringLiteral("label_" + i); + } + FunctionCallExpr overlap = new FunctionCallExpr("arrays_overlap", Arrays.asList( + new SlotRef(null, "labels"), new ArrayLiteral(ArrayType.create(Type.STRING, true), labels))); + overlap.setFn(new ScalarFunction(new FunctionName("arrays_overlap"), + Arrays.asList(ArrayType.create(Type.STRING, true), ArrayType.create(Type.STRING, true)), + Type.BOOLEAN, false, true)); + LancePredicateConverter.ConversionResult converted = new LancePredicateConverter(schema) + .convert(Collections.singletonList(overlap)); + Assertions.assertTrue(converted.getResidualConjuncts().isEmpty()); + LanceScalarIndexPlanner.Plan selected = LanceScalarIndexPlanner.plan(metadata, + converted.getPushedConjuncts(), fragments); + if (count == 64) { + Assertions.assertNotNull(selected); + Assertions.assertEquals(1, selected.splits.splitCount()); + // Count the enclosing AND and its other leaf, not just overlap's 127 nodes. + Assertions.assertNull(LanceScalarIndexPlanner.plan(metadata, + Arrays.asList(overlap, contains("extra")), fragments)); + } else { + Assertions.assertNull(selected); + } + } + } + + @Test + public void testWideConjunctionUsesBalancedDepth() { + List filters = new ArrayList<>(); + filters.add(equal("key", 1)); + for (int i = 0; i < 33; i++) { + filters.add(equal("filter_" + i, i)); + } + Assertions.assertNotNull(plan(filters, false)); + } + + @Test + public void testShortInExpansionAtOverlapBudgetBoundary() throws Exception { + Schema schema = new Schema(Arrays.asList(new Field("labels", + FieldType.nullable(ArrowType.List.INSTANCE), + Collections.singletonList(Field.nullable("item", ArrowType.Utf8.INSTANCE))), + Field.nullable("category", new ArrowType.Int(64, true)))); + List fragments = Arrays.asList( + new LanceFragmentInfo(1, 10, 10), new LanceFragmentInfo(2, 10, 10)); + Map fields = new HashMap<>(); + fields.put("labels", 9); + fields.put("category", 10); + LanceTableMetadata metadata = LanceTableMetadata.createSnapshotWithIndexes( + new LanceTableAccess("s3://bucket/labels.lance", Collections.emptyMap()), 42, schema, + fragments, fields, Arrays.asList( + new LanceIndexSegmentInfo(UUID.randomUUID(), "labels_idx", Collections.singletonList(9), + Arrays.asList(1L, 2L), IndexType.LABEL_LIST, null), + new LanceIndexSegmentInfo(UUID.randomUUID(), "category_idx", Collections.singletonList(10), + Arrays.asList(1L, 2L), IndexType.BTREE, null))); + Map visible = new HashMap<>(); + fragments.forEach(fragment -> visible.put(fragment.getId(), fragment)); + for (int count : Arrays.asList(62, 63)) { + StringLiteral[] labels = new StringLiteral[count]; + for (int i = 0; i < count; i++) { + labels[i] = new StringLiteral("label_" + i); + } + Expr overlap = new FunctionCallExpr("arrays_overlap", Arrays.asList( + new SlotRef(null, "labels"), new ArrayLiteral(ArrayType.create(Type.STRING, true), labels))); + Expr in = new InPredicate(new SlotRef(null, "category"), + Arrays.asList(new IntLiteral(0), new IntLiteral(1)), false); + LanceScalarIndexPlanner.Plan selected = LanceScalarIndexPlanner.plan(metadata, + Arrays.asList(overlap, in), visible); + if (count == 62) { + Assertions.assertNotNull(selected); + } else { + Assertions.assertNull(selected); + } + } + } + + @Test + public void testLiteralAndRefinedPrefixes() { + Expr key = new BinaryPredicate(BinaryPredicate.Operator.EQ, + new SlotRef(null, "key"), new StringLiteral("x")); + for (String prefix : Arrays.asList("test_ns$", "literal%value", "literal\\value")) { + Expr startsWith = new FunctionCallExpr("starts_with", + Arrays.asList(new SlotRef(null, "key"), new StringLiteral(prefix))); + Assertions.assertNotNull(plan(Collections.singletonList(startsWith), true)); + Assertions.assertNotNull(plan(Collections.singletonList(or(key, startsWith)), true)); + } + for (String pattern : Arrays.asList("foo%bar%", "foo_bar%")) { + Expr like = new LikePredicate(LikePredicate.Operator.LIKE, + new SlotRef(null, "key"), new StringLiteral(pattern)); + Assertions.assertNotNull(plan(Collections.singletonList(like), true)); + // Native cannot union an index branch that still carries a refine predicate. + Assertions.assertNull(plan(Collections.singletonList(or(key, like)), true)); + Assertions.assertNull(plan(Collections.singletonList(or(and(key, like), key)), true)); + } + } + + private static FunctionCallExpr contains(String label) { + FunctionCallExpr function = new FunctionCallExpr("array_contains", + Arrays.asList(new SlotRef(null, "labels"), new StringLiteral(label))); + function.setFn(new ScalarFunction(new FunctionName("array_contains"), + Arrays.asList(ArrayType.create(Type.STRING, true), Type.STRING), Type.BOOLEAN, false, true)); + return function; + } + + private static LanceScalarIndexPlanner.Plan plan(Expr filter) { + return plan(Collections.singletonList(filter), false); + } + + private static LanceScalarIndexPlanner.Plan plan(List filters, boolean stringKey) { + LanceFragmentInfo first = new LanceFragmentInfo(1, 10, 10); + LanceFragmentInfo second = new LanceFragmentInfo(2, 10, 10); + Map fragments = new HashMap<>(); + fragments.put(1L, first); + fragments.put(2L, second); + Map fields = new HashMap<>(); + fields.put("key", 9); + fields.put("other", 10); + LanceTableMetadata metadata = LanceTableMetadata.createSnapshotWithIndexes( + new LanceTableAccess("s3://bucket/labels.lance", Collections.emptyMap()), 42, + new Schema(Arrays.asList(Field.nullable("key", stringKey ? ArrowType.Utf8.INSTANCE : new ArrowType.Int(64, true)), + Field.nullable("other", new ArrowType.Int(64, true)))), + Arrays.asList(first, second), fields, Collections.singletonList(new LanceIndexSegmentInfo( + UUID.randomUUID(), "key_idx", Collections.singletonList(9), + Arrays.asList(1L, 2L), IndexType.BTREE, null))); + return LanceScalarIndexPlanner.plan(metadata, filters, fragments); + } + + private static Expr equal(String column, int value) { + return new BinaryPredicate(BinaryPredicate.Operator.EQ, new SlotRef(null, column), new IntLiteral(value)); + } + + private static Expr and(Expr left, Expr right) { + return new CompoundPredicate(CompoundPredicate.Operator.AND, left, right); + } + + private static Expr or(Expr left, Expr right) { + return new CompoundPredicate(CompoundPredicate.Operator.OR, left, right); + } + + private static Expr not(Expr child) { + return new CompoundPredicate(CompoundPredicate.Operator.NOT, child, null); + } +} diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/source/LanceScanNodeTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/source/LanceScanNodeTest.java index 6131a2fcb2c897..7e58ad90413445 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/source/LanceScanNodeTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/source/LanceScanNodeTest.java @@ -234,7 +234,7 @@ public void testScalarSegmentPlanRejectsUnsafeCoverageAndUnsupportedTypes() thro } @Test - public void testScalarSegmentSelectionDoesNotDescendThroughOrOrNot() throws Exception { + public void testScalarSegmentSelectionSupportsSameColumnOrAndNot() throws Exception { Expr predicate = scalarPredicate(); for (Expr filter : Arrays.asList( new CompoundPredicate(CompoundPredicate.Operator.OR, predicate, predicate), @@ -243,11 +243,29 @@ public void testScalarSegmentSelectionDoesNotDescendThroughOrOrNot() throws Exce setMetadata(node, scalarMetadata(Collections.singletonList( scalarSegment(UUID.randomUUID(), IndexType.BTREE, Arrays.asList(1L, 2L, 3L, 4L))))); setPushedConjuncts(node, filter); - Assert.assertEquals(4, node.getSplits(20).size()); + int expectedSplits = ((CompoundPredicate) filter).getOp() == CompoundPredicate.Operator.NOT ? 4 : 1; + List initialSplits = node.getSplits(20); + Assert.assertEquals(expectedSplits, initialSplits.size()); + if (expectedSplits == 4) { + for (Split split : initialSplits) { + TFileRangeDesc range = new TFileRangeDesc(); + node.setScanParams(range, split); + TLanceFileDesc params = range.getTableFormatParams().getLanceParams(); + Assert.assertFalse(params.isSetIndexSegmentUuids()); + Assert.assertTrue(params.isSetUseScalarIndex()); + Assert.assertFalse(params.isUseScalarIndex()); + } + } setPushedConjuncts(node, new CompoundPredicate(CompoundPredicate.Operator.AND, filter, predicate)); List splits = node.getSplits(20); Assert.assertEquals(1, splits.size()); - Assert.assertTrue(((LanceSplit) splits.get(0)).getIndexSegmentUuid().isPresent()); + List assignedFragments = new ArrayList<>(); + for (Split split : splits) { + assignedFragments.addAll(((LanceSplit) split).getFragmentIds()); + Assert.assertTrue(((LanceSplit) split).getIndexSegmentUuid().isPresent()); + } + Collections.sort(assignedFragments); + Assert.assertEquals(Arrays.asList(1L, 2L, 3L, 4L), assignedFragments); } } @@ -286,6 +304,51 @@ public void testScalarSegmentDoesNotPushLimitPastDorisResidual() throws Exceptio Assert.assertFalse(range.getTableFormatParams().getLanceParams().isSetLimit()); } + @Test + public void testMissingSegmentMetadataPreservesNativeIndexSelection() throws Exception { + List segments = Collections.singletonList( + scalarSegment(UUID.randomUUID(), IndexType.BTREE, Arrays.asList(1L, 2L, 3L, 4L))); + LanceTableMetadata complete = scalarMetadata(segments); + LanceTableAccess access = new LanceTableAccess("s3://bucket/scalar.lance", Collections.emptyMap()); + for (LanceTableMetadata metadata : Arrays.asList( + LanceTableMetadata.createSnapshotWithUnavailableFieldIds(access, 42, + complete.getSchema(), complete.getFragments(), segments), + LanceTableMetadata.createSnapshotWithIndexes(access, 42, + complete.getSchema(), complete.getFragments(), Collections.singletonMap("key", 9), + Collections.emptyList()))) { + assertNativeFragmentSelection(metadata); + } + } + + @Test + public void testCompetingLogicalIndexesPreserveNativeSelection() throws Exception { + List segments = Arrays.asList( + new LanceIndexSegmentInfo(UUID.randomUUID(), "a_key_idx", Collections.singletonList(9), + Collections.singletonList(1L), IndexType.BTREE, null), + new LanceIndexSegmentInfo(UUID.randomUUID(), "z_key_idx", Collections.singletonList(9), + Arrays.asList(1L, 2L, 3L, 4L), IndexType.BTREE, null)); + assertNativeFragmentSelection(scalarMetadata(segments)); + Collections.reverse(segments); + assertNativeFragmentSelection(scalarMetadata(segments)); + } + + private static void assertNativeFragmentSelection(LanceTableMetadata metadata) throws Exception { + LanceScanNode node = newNode(); + setMetadata(node, metadata); + setPushedConjuncts(node, new CompoundPredicate(CompoundPredicate.Operator.OR, + new BinaryPredicate(BinaryPredicate.Operator.EQ, new SlotRef(null, "key"), new IntLiteral(1)), + new BinaryPredicate(BinaryPredicate.Operator.EQ, new SlotRef(null, "key"), new IntLiteral(2)))); + List splits = node.getSplits(20); + Assert.assertEquals(4, splits.size()); + for (Split split : splits) { + TFileRangeDesc range = new TFileRangeDesc(); + node.setScanParams(range, split); + TLanceFileDesc params = range.getTableFormatParams().getLanceParams(); + Assert.assertFalse(params.isSetIndexSegmentUuids()); + Assert.assertFalse(params.isSetUseScalarIndex()); + } + } + private static Expr scalarPredicate() { return new BinaryPredicate(BinaryPredicate.Operator.GE, new SlotRef(null, "key"), new IntLiteral(2)); } diff --git a/regression-test/suites/external_table_p0/lance/test_lance_array_predicate_pushdown.groovy b/regression-test/suites/external_table_p0/lance/test_lance_array_predicate_pushdown.groovy new file mode 100644 index 00000000000000..7c2d21a8e029ef --- /dev/null +++ b/regression-test/suites/external_table_p0/lance/test_lance_array_predicate_pushdown.groovy @@ -0,0 +1,154 @@ +// 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')" + String green = "array_contains(labels, 'green')" + def cases = [ + [red, [0, 2, 5, 7, 8, 10, 13, 15], 8, 4], + ["${red} AND ${blue}", [2, 7, 10, 15], 4, 2], + ["${red} OR ${blue}", [0, 1, 2, 5, 6, 7, 8, 9, 10, 13, 14, 15], 12, 6], + // Nereids rewrites this shape to arrays_overlap before Lance conversion. + ["${red} OR ${blue} OR ${green}", [0, 1, 2, 5, 6, 7, 8, 9, 10, 13, 14, 15], 12, 6], + ["arrays_overlap(labels, ['red', 'blue', 'green'])", + [0, 1, 2, 5, 6, 7, 8, 9, 10, 13, 14, 15], 12, 6], + ["NOT arrays_overlap(labels, ['red', 'blue', 'green'])", [3, 11], -1, -1], + ["NOT (${red})", [1, 3, 6, 9, 11, 14], -1, -1], + ["NOT (${green})", [0, 1, 2, 3, 5, 6, 7, 8, 9, 10, 11, 13, 14, 15], -1, -1], + ["category IN (0, 1)", [0, 1, 3, 4, 6, 7, 9, 10, 12, 13, 15], 11, 6], + ["category = 0 OR category = 1", [0, 1, 3, 4, 6, 7, 9, 10, 12, 13, 15], 11, 6], + ["${red} AND category = 1", [7, 10, 13], null], + ["${red} AND category <> 1", [0, 2, 5, 8, 15], null] + ] + for (int size : [64, 65]) { + String needles = (["'red'"] + (1.. + 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" && c[2] != -1) { + contains "lanceScalarIndexScan=SEGMENT" + } else { + notContains "lanceScalarIndexScan=SEGMENT" + contains "lanceFragmentGrouping=FRAGMENT" + contains "lanceFragments=2" + } + } + 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 != "unindexed") { + String profile = profiles.getProfileBySql(token, + ["LanceScalarIndexSegmentsSearched", "LanceScalarIndexCandidateRows", + "LanceScalarIndexSegmentFallbacks"]) + def counter = { String name -> + def matches = profile =~ /${name}: (?:sum )?(\d+)\b/ + assertTrue(matches.find(), "Missing counter ${name}: ${profile}") + return matches.group(1).toLong() + } + // Broad complements stay as parallel scans without repeated segment searches. + long expectedSearches = c[2] == -1 ? 0L : 1L + assertEquals(expectedSearches, counter("LanceScalarIndexSegmentsSearched")) + long expectedCandidates = c[2] == null ? -1L : + (c[table == "partial" ? 3 : 2] as Number).longValue() + if (expectedCandidates >= 0) { + assertEquals(expectedCandidates, 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]], + ["arrays_overlap(labels, ['green', NULL])", [6, 14]], + ["arrays_overlap(labels, [])", []], + ["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() }) + // An unordered LIMIT reaches the scanner; TopN's bound does not. The two + // selective predicates force results from the indexed and uncovered domains. + for (int expectedId : [5, 13]) { + String limited = "SELECT id FROM ${relation} WHERE ${red} AND id = ${expectedId} LIMIT 1" + explain { + sql(limited) + contains "lanceLimit=1" + contains "lancePushdownPredicate=" + notContains "predicates:" + } + assertEquals([[(long) expectedId]], sql(limited)) + } + // 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")) + } + } finally { + sql "DROP CATALOG IF EXISTS ${catalog}" + } +} diff --git a/thirdparty/build-thirdparty.sh b/thirdparty/build-thirdparty.sh index a5e209a0aedb32..89c69090e04c49 100755 --- a/thirdparty/build-thirdparty.sh +++ b/thirdparty/build-thirdparty.sh @@ -156,6 +156,8 @@ if [[ ! -f "${TP_DIR}/vars.sh" ]]; then fi . "${TP_DIR}/vars.sh" +. "${TP_DIR}/lance-install.sh" +LANCE_C_INSTALL_FINGERPRINT="$(lance_c_install_fingerprint "${TP_DIR}")" cd "${TP_DIR}" @@ -2228,9 +2230,14 @@ build_lance_c() { env "${cargo_env[@]}" "${cargo_bin}" "${cargo_args[@]}" mkdir -p "${TP_INSTALL_DIR}/include" "${TP_INSTALL_DIR}/lib64" + # Invalidate before publishing either file so interrupted installs cannot reuse + # a matching marker with a partial header/archive pair. + rm -f "${TP_INSTALL_DIR}/lib64/.lance-c-fingerprint" rm -rf "${TP_INSTALL_DIR}/include/lance" cp -av include/lance "${TP_INSTALL_DIR}/include/" install_rust_archive "${BUILD_DIR}/release/liblance_c.a" + printf '%s\n' "${LANCE_C_INSTALL_FINGERPRINT}" > "${TP_INSTALL_DIR}/lib64/.lance-c-fingerprint.tmp" + mv "${TP_INSTALL_DIR}/lib64/.lance-c-fingerprint.tmp" "${TP_INSTALL_DIR}/lib64/.lance-c-fingerprint" } # paimon-rust diff --git a/thirdparty/lance-install.sh b/thirdparty/lance-install.sh new file mode 100644 index 00000000000000..28594d61c79026 --- /dev/null +++ b/thirdparty/lance-install.sh @@ -0,0 +1,38 @@ +#!/usr/bin/env bash +# 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. + +# Keep the check independent of the installed prefix: an external compilation image +# can carry old vars.sh alongside an ABI-compatible but behaviorally stale archive. +lance_c_install_fingerprint() ( + local definitions="$1" + local TP_DIR="${definitions}" + source "${definitions}/vars.sh" || return 1 + local patch_checksum + patch_checksum="$(cksum < "${definitions}/patches/lance-c-foyer.patch")" || return 1 + printf '%s\n' "${LANCE_C_SOURCE}" "${LANCE_C_MD5SUM}" "${patch_checksum}" +) + +lance_c_install_is_current() { + local definitions="$1" installed="$2" expected + [[ -s "${installed}/lib64/liblance_c.a" && + -s "${installed}/include/lance/lance.h" && + -s "${installed}/include/lance/lance.hpp" && + -s "${installed}/lib64/.lance-c-fingerprint" ]] || return 1 + expected="$(lance_c_install_fingerprint "${definitions}")" || return 1 + [[ "$(cat "${installed}/lib64/.lance-c-fingerprint")" == "${expected}" ]] +} diff --git a/thirdparty/patches/lance-c-foyer.patch b/thirdparty/patches/lance-c-foyer.patch index 73f244c49fd57d..bb2104ab32d66d 100644 --- a/thirdparty/patches/lance-c-foyer.patch +++ b/thirdparty/patches/lance-c-foyer.patch @@ -1,8 +1,7 @@ # Foyer data-cache integration for lance-format/lance-c#73. -# Base: 9bd730add2ac70316c1d642b8459011e2dd92022 +# Base: cd63420bfbe27f6f0a1edcc873b9191af7d52852 (lance-format/lance-c main, including #93) # Source: https://github.com/Gabriel39/lance-c/commit/24c7ca4bcb9422c113b0d3e07e4efe1173b0bc9f -# Regenerate the payload in lance-c; do not maintain separate downstream edits: -# git diff --full-index --binary 9bd730add2ac70316c1d642b8459011e2dd92022 24c7ca4bcb9422c113b0d3e07e4efe1173b0bc9f -- . | sed 's/^ $//' +# The unchanged Foyer payload is reapplied to the newer upstream base with git apply --3way. diff --git a/Cargo.lock b/Cargo.lock index 78812f1330e146db295a14276f90f9e28654f399..daf9f85daf3e0ea59bb906e8e3c32470519b084b 100644 --- a/Cargo.lock @@ -199,7 +198,7 @@ index 18654b862440a90aa6bf0dd846e8f1ad6036fd7c..fd134a7151e1df61da201d867a642e9b prost = "0.14" snafu = "0.9" diff --git a/README.md b/README.md -index 9ceccdc4f3f6b5d13dc271be3d2659c7be400306..94de4d1c26d8ab869d1726e2b416d21f21661d6d 100644 +index 9f0871b17b13ebfd96916947b4bee0dd06f6a468..441719c9f3904058127e55d1cd29ea1afc878935 100644 --- a/README.md +++ b/README.md @@ -68,6 +68,7 @@ Based on the [liblance RFC](https://github.com/lance-format/lance/discussions/60 @@ -208,9 +207,9 @@ index 9ceccdc4f3f6b5d13dc271be3d2659c7be400306..94de4d1c26d8ab869d1726e2b416d21f | [x] | Filter pushdown | `lance_scanner_set_substrait_filter()` accepts a serialized Substrait `ExtendedExpression`; `lance_scanner_additional_sql_filter()` adds SQL predicates with AND before scanning starts | +| [x] | Data-file cache | Optional Foyer memory/disk cache for immutable `data/*.lance` reads | - ## Multi-vector search + ## Segment-scoped array label filters -@@ -231,6 +232,35 @@ auto ds = lance::Dataset::open_with_session(session, "data.lance"); +@@ -362,6 +363,35 @@ auto ds = lance::Dataset::open_with_session(session, "data.lance"); auto stats = session.cache_stats(); ``` @@ -247,7 +246,7 @@ index 9ceccdc4f3f6b5d13dc271be3d2659c7be400306..94de4d1c26d8ab869d1726e2b416d21f `lance_dataset_open` takes a `version` argument — `0` means the latest, any diff --git a/include/lance/lance.h b/include/lance/lance.h -index 7630913fd7849d0bff2cb0f31e1ab1199cb9f0fb..1c1d0752efdda5023b29eafe5aa9a87e03dc5303 100644 +index 94772e2a14bf9a8a950dfb2d39112136e27197c0..fe48eb0d9afd1c7157a3bd2230cce7e769120bac 100644 --- a/include/lance/lance.h +++ b/include/lance/lance.h @@ -214,6 +214,36 @@ typedef struct LanceSessionCacheStats { @@ -335,7 +334,7 @@ index 7630913fd7849d0bff2cb0f31e1ab1199cb9f0fb..1c1d0752efdda5023b29eafe5aa9a87e void lance_dataset_close(LanceDataset* dataset); diff --git a/include/lance/lance.hpp b/include/lance/lance.hpp -index 286724e96736e354785e046f2fcfe5ae7c65147b..5070d423453a5d60d31d03ce758ad7dbe83a6ea2 100644 +index edcd39fc3caba387ab490fa84139f24e9376130a..17b71957d0f5c419fe0b3392ce11faa757ddbfe1 100644 --- a/include/lance/lance.hpp +++ b/include/lance/lance.hpp @@ -176,6 +176,13 @@ struct SqlColumn { @@ -2244,7 +2243,7 @@ index 1971510de4a13ee6f04fc39c159b088e78e21d55..ba51c87a8cebb6897c48fb55509f71da // SAFETY: `out_dataset` is non-NULL (checked above) and the caller // guarantees it points to caller-owned, writable storage of size diff --git a/tests/c_api_test.rs b/tests/c_api_test.rs -index 1a25e34a8ca50623812ffedaa7bb5a8c54f7850e..c341859e9434d4578b2833eca6ec313a28f31f32 100644 +index d3bea44529a842fe3ae28ee3f90cd8d45fe0b394..98e6901a2cfeb2cc98df6606e26b04788d4218b2 100644 --- a/tests/c_api_test.rs +++ b/tests/c_api_test.rs @@ -100,10 +100,83 @@ fn create_large_dataset(num_rows: i32) -> (tempfile::TempDir, String) { @@ -2503,7 +2502,7 @@ index 1a25e34a8ca50623812ffedaa7bb5a8c54f7850e..c341859e9434d4578b2833eca6ec313a fn test_dataset_restore_to_current_latest_writes_new_manifest() { // Restoring to the current latest still writes a new manifest. The diff --git a/tests/cpp/test_c_api.c b/tests/cpp/test_c_api.c -index 5df9cde4d8fa971823b9e6c974a8a0db6c9ec9ca..1444e55161631ac8ad976f427d867ad91d358195 100644 +index efad5ed55987976a3f1820220e42fc58bef42638..050543aebc02454cab09f6e0bd38ec10b1d21720 100644 --- a/tests/cpp/test_c_api.c +++ b/tests/cpp/test_c_api.c @@ -171,6 +171,37 @@ static void test_shared_session(const char *uri) { @@ -2544,16 +2543,16 @@ index 5df9cde4d8fa971823b9e6c974a8a0db6c9ec9ca..1444e55161631ac8ad976f427d867ad9 static void test_scan(const char *uri) { printf(" test_scan... "); -@@ -1323,6 +1354,7 @@ int main(int argc, char **argv) { - +@@ -1471,6 +1502,7 @@ int main(int argc, char **argv) { + test_batch_nearest(uri); test_open_and_metadata(uri); test_shared_session(uri); + test_data_cache_session(uri, write_uri); test_scan(uri); + test_distance_range(uri); test_scan_with_limit(uri); - test_scanner_blob_handling(blob_uri); diff --git a/tests/cpp/test_cpp_api.cpp b/tests/cpp/test_cpp_api.cpp -index 0332d1a2687354984d0551e10716df08344a4013..ef64ab0d02bed2f1750c3d388792e4aa44265e24 100644 +index e2d541e4430889aba89ab85c01cc7e3eacf97baa..ac1c10bccfcf399bbbfcfe908d3188484aff0add 100644 --- a/tests/cpp/test_cpp_api.cpp +++ b/tests/cpp/test_cpp_api.cpp @@ -131,6 +131,28 @@ static void test_shared_session(const std::string& uri) { @@ -2585,11 +2584,11 @@ index 0332d1a2687354984d0551e10716df08344a4013..ef64ab0d02bed2f1750c3d388792e4aa static void test_dataset_schema(const std::string& uri) { TEST(test_dataset_schema); -@@ -1211,6 +1233,7 @@ int main(int argc, char** argv) { +@@ -1353,6 +1375,7 @@ int main(int argc, char** argv) { test_dataset_open(uri); test_shared_session(uri); + test_data_cache_session(uri, write_uri); test_dataset_schema(uri); test_scanner_fluent(uri); - test_scanner_async_stream_ownership(uri); + test_distance_range(uri); diff --git a/thirdparty/test/lance-install-test.sh b/thirdparty/test/lance-install-test.sh new file mode 100644 index 00000000000000..db93c355c0f421 --- /dev/null +++ b/thirdparty/test/lance-install-test.sh @@ -0,0 +1,176 @@ +#!/usr/bin/env bash +# 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. + +set -euo pipefail +ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/../.." && pwd)" +work="$(mktemp -d)" +trap 'rm -rf "${work}"' EXIT +mkdir -p "${work}/repo/thirdparty/patches" "${work}/external" +cp "${ROOT}/thirdparty/vars.sh" "${work}/repo/thirdparty/" +cp "${ROOT}/thirdparty/patches/lance-c-foyer.patch" "${work}/repo/thirdparty/patches/" +# Extract the real build gate; all destructive operations stay inside this temporary install. +echo 'set -eo pipefail' > "${work}/gate.sh" +sed -n '/^# build thirdparty libraries if necessary/,/^update_submodule()/p' "${ROOT}/build.sh" \ + | sed '$d' >> "${work}/gate.sh" +export DORIS_HOME="${work}/repo" DORIS_THIRDPARTY="${work}/external" +export TARGET_SYSTEM=Linux CLEAN=0 PARALLEL=1 +export TEST_HELPER="${ROOT}/thirdparty/lance-install.sh" +cp "${TEST_HELPER}" "${work}/repo/thirdparty/" +cp -r "${DORIS_HOME}/thirdparty/." "${DORIS_THIRDPARTY}/" +printf '#!/usr/bin/env bash\nexit 0\n' > "${DORIS_THIRDPARTY}/download-thirdparty.sh" +cat > "${DORIS_THIRDPARTY}/build-thirdparty.sh" <<'BUILDER' +set -euo pipefail +TP_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +source "${TP_DIR}/lance-install.sh" +fingerprint="$(lance_c_install_fingerprint "${TP_DIR}")" +echo rebuilt >> "${DORIS_THIRDPARTY}/builds" +installed="${DORIS_THIRDPARTY}/installed" +mkdir -p "${installed}/lib/hadoop_hdfs/native" "${installed}/lib/hadoop_hdfs_3_4/native" "${installed}/lib64" \ + "${installed}/include/lance" "${installed}/include/paimon_rust" +for file in lib/hadoop_hdfs/native/libhdfs.a lib/hadoop_hdfs_3_4/native/libhdfs.a \ + lib64/liblance_c.a lib64/libpaimon_c.a \ + include/lance/lance.h include/lance/lance.hpp include/paimon_rust/paimon.h; do + echo artifact > "${installed}/${file}" +done +printf '%s\n' "${fingerprint}" > "${installed}/lib64/.lance-c-fingerprint" +BUILDER +# A complete legacy image has all old sentinels but no Lance revision marker. +installed="${DORIS_THIRDPARTY}/installed" +mkdir -p "${installed}/lib/hadoop_hdfs/native" "${installed}/lib/hadoop_hdfs_3_4/native" "${installed}/lib64" \ + "${installed}/include/lance" "${installed}/include/paimon_rust" +for file in lib/hadoop_hdfs/native/libhdfs.a lib/hadoop_hdfs_3_4/native/libhdfs.a \ + lib64/liblance_c.a lib64/libpaimon_c.a \ + include/lance/lance.h include/lance/lance.hpp include/paimon_rust/paimon.h; do + echo legacy > "${installed}/${file}" +done +bash "${work}/gate.sh" +[[ -f "${DORIS_THIRDPARTY}/builds" ]] || { echo 'FAIL: reused unversioned Lance archive'; exit 1; } +source "${TEST_HELPER}" +lance_c_install_is_current "${DORIS_HOME}/thirdparty" "${installed}" +check_builds() { + [[ "$(wc -l < "${DORIS_THIRDPARTY}/builds")" -eq "$1" ]] +} +check_builds 1 +bash "${work}/gate.sh" +check_builds 1 +echo 'PASS: legacy image rebuilt; matching install reused' +# Failed preflight must preserve every dependency, not just the Lance archive. +expect_preserved_install() { + cp -a "${installed}" "${work}/saved-install" + cp "${DORIS_THIRDPARTY}/builds" "${work}/saved-builds" + if bash "${work}/gate.sh" > "${work}/preflight.log" 2>&1; then + echo 'FAIL: accepted incomplete or mismatched rebuild sources'; exit 1 + fi + if ! diff -r "${work}/saved-install" "${installed}"; then + echo 'FAIL: rejected rebuild sources after changing installed dependencies'; exit 1 + fi + cmp "${work}/saved-builds" "${DORIS_THIRDPARTY}/builds" + rm -rf "${work}/saved-install" +} +echo stale > "${installed}/lib64/.lance-c-fingerprint" +for file in lance-install.sh vars.sh patches/lance-c-foyer.patch download-thirdparty.sh build-thirdparty.sh; do + mv "${DORIS_THIRDPARTY}/${file}" "${work}/missing-input" + expect_preserved_install + mv "${work}/missing-input" "${DORIS_THIRDPARTY}/${file}" +done +echo 'PASS: incomplete external sources rejected before changing installed dependencies' +echo stale > "${installed}/lib64/.lance-c-fingerprint" +bash "${work}/gate.sh" +check_builds 2 +# Pin and patch updates both invalidate installed artifacts, even with identical ABI. +sed 's/^LANCE_C_SOURCE=.*/LANCE_C_SOURCE="lance-c-test-revision"/' "${DORIS_HOME}/thirdparty/vars.sh" \ + > "${work}/new-vars.sh" +mv "${work}/new-vars.sh" "${DORIS_HOME}/thirdparty/vars.sh" +expect_preserved_install +cp "${DORIS_HOME}/thirdparty/vars.sh" "${DORIS_THIRDPARTY}/vars.sh" +bash "${work}/gate.sh" +check_builds 3 +echo '# test patch update' >> "${DORIS_HOME}/thirdparty/patches/lance-c-foyer.patch" +expect_preserved_install +cp "${DORIS_HOME}/thirdparty/patches/lance-c-foyer.patch" "${DORIS_THIRDPARTY}/patches/" +bash "${work}/gate.sh" +check_builds 4 +echo 'PASS: mismatched sources rejected; synchronized pin and patch updates rebuild' +rm "${installed}/include/lance/lance.h" +bash "${work}/gate.sh" +check_builds 5 +: > "${installed}/lib64/liblance_c.a" +bash "${work}/gate.sh" +check_builds 6 +echo 'PASS: incomplete header/archive install rebuilt' +# A legacy external builder may exit successfully without installing the new revision. +cp "${DORIS_THIRDPARTY}/build-thirdparty.sh" "${work}/good-builder.sh" +printf '#!/usr/bin/env bash\nexit 0\n' > "${DORIS_THIRDPARTY}/build-thirdparty.sh" +rm "${installed}/lib64/.lance-c-fingerprint" +if bash "${work}/gate.sh" > "${work}/old-builder.log" 2>&1; then + echo 'FAIL: accepted output from a stale external builder'; exit 1 +fi +grep -q 'Lance dependency revision does not match' "${work}/old-builder.log" +cp "${work}/good-builder.sh" "${DORIS_THIRDPARTY}/build-thirdparty.sh" +bash "${DORIS_THIRDPARTY}/build-thirdparty.sh" +echo 'PASS: stale external builder cannot silently satisfy the revision gate' +# Images without rebuild sources must fail before deleting installed dependencies. +rm "${DORIS_THIRDPARTY}/build-thirdparty.sh" "${installed}/lib64/.lance-c-fingerprint" +if bash "${work}/gate.sh" > "${work}/missing-source.log" 2>&1; then + echo 'FAIL: accepted stale compilation image without build sources'; exit 1 +fi +[[ -s "${installed}/lib64/liblance_c.a" ]] +echo 'PASS: missing rebuild source fails without deleting installed artifacts' +# Execute the actual publication function with a tiny Cargo stand-in. A failed +# archive copy must invalidate the old marker, and only a complete retry may stamp it. +sed -n '/^build_lance_c()/,/^}/p' "${ROOT}/thirdparty/build-thirdparty.sh" > "${work}/publish-function.sh" +export TP_DIR="${DORIS_HOME}/thirdparty" +source "${TP_DIR}/vars.sh" +export LANCE_C_SOURCE TP_SOURCE_DIR TP_INSTALL_DIR +export LANCE_C_INSTALL_FINGERPRINT="$(lance_c_install_fingerprint "${TP_DIR}")" +mkdir -p "${TP_SOURCE_DIR}/${LANCE_C_SOURCE}/include/lance" "${TP_INSTALL_DIR}/bin" +echo header > "${TP_SOURCE_DIR}/${LANCE_C_SOURCE}/include/lance/lance.h" +echo header > "${TP_SOURCE_DIR}/${LANCE_C_SOURCE}/include/lance/lance.hpp" +printf '#!/bin/sh\nexit 0\n' > "${TP_INSTALL_DIR}/bin/protoc" +chmod +x "${TP_INSTALL_DIR}/bin/protoc" +cat > "${work}/cargo" <<'CARGO' +#!/usr/bin/env bash +set -eu +if [[ "$1" == --version ]]; then + echo 'cargo 1.94.0'; exit 0 +fi +mkdir -p "${CARGO_TARGET_DIR}/release" +echo archive > "${CARGO_TARGET_DIR}/release/liblance_c.a" +CARGO +chmod +x "${work}/cargo" +export LANCE_C_CARGO="${work}/cargo" RUSTUP_TOOLCHAIN=1.94.0 +export BUILD_DIR=build KERNEL=Linux LANCE_C_CARGO_OFFLINE=OFF +cat > "${work}/publish.sh" <<'PUBLISH' +set -eo pipefail +check_if_source_exist() { :; } +install_rust_archive() { + if [[ "${FAIL_INSTALL:-0}" == 1 ]]; then return 1; fi + cp "$1" "${TP_INSTALL_DIR}/lib64/liblance_c.a" +} +source "$1" +build_lance_c +PUBLISH +bash "${work}/publish.sh" "${work}/publish-function.sh" > "${work}/publish.log" 2>&1 +lance_c_install_is_current "${TP_DIR}" "${TP_INSTALL_DIR}" +if FAIL_INSTALL=1 bash "${work}/publish.sh" "${work}/publish-function.sh" >> "${work}/publish.log" 2>&1; then + echo 'FAIL: expected archive publication failure'; exit 1 +fi +[[ ! -e "${TP_INSTALL_DIR}/lib64/.lance-c-fingerprint" ]] +bash "${work}/publish.sh" "${work}/publish-function.sh" >> "${work}/publish.log" 2>&1 +lance_c_install_is_current "${TP_DIR}" "${TP_INSTALL_DIR}" +echo 'PASS: successful install stamped; failed publication invalidated; retry repaired' diff --git a/thirdparty/vars.sh b/thirdparty/vars.sh index fc5e9818c32dd4..147e4e4b2c45d6 100644 --- a/thirdparty/vars.sh +++ b/thirdparty/vars.sh @@ -580,11 +580,11 @@ PUGIXML_SOURCE=pugixml-1.15 PUGIXML_MD5SUM="3b894c29455eb33a40b165c6e2de5895" # lance-c -# Complete-segment prefilter fixes are supplied by upstream lance-c, not local patches. -LANCE_C_DOWNLOAD="https://codeload.github.com/lance-format/lance-c/tar.gz/9bd730add2ac70316c1d642b8459011e2dd92022" -LANCE_C_NAME="lance-c-9bd730add2ac70316c1d642b8459011e2dd92022.tar.gz" -LANCE_C_SOURCE="lance-c-9bd730add2ac70316c1d642b8459011e2dd92022" -LANCE_C_MD5SUM="63851b09bf1689032579f1a094ff2f37" +# Includes lance-c #93: scoped boolean scalar-index expressions and Substrait label filters. +LANCE_C_DOWNLOAD="https://codeload.github.com/lance-format/lance-c/tar.gz/cd63420bfbe27f6f0a1edcc873b9191af7d52852" +LANCE_C_NAME="lance-c-cd63420bfbe27f6f0a1edcc873b9191af7d52852.tar.gz" +LANCE_C_SOURCE="lance-c-cd63420bfbe27f6f0a1edcc873b9191af7d52852" +LANCE_C_MD5SUM="37d82907559c4fdb8ed6e5b767b0d309" # paimon-rust PAIMON_RUST_DOWNLOAD="https://github.com/apache/paimon-rust/archive/refs/tags/v0.4.0-rc1.tar.gz"