diff --git a/be/src/format_v2/table/lance_reader.cpp b/be/src/format_v2/table/lance_reader.cpp index 4f1229901a9a49..142c35ac4d42bb 100644 --- a/be/src/format_v2/table/lance_reader.cpp +++ b/be/src/format_v2/table/lance_reader.cpp @@ -683,6 +683,14 @@ Status LanceTableReader::_open_scanner(const TFileRangeDesc& range) { return _lance_error("set Lance scanner fragment ids"); } } + // Ordinary scans may carry a pushed-down LIMIT. The FE only sets it when all predicates are + // pushed into Lance, so the scanner can safely stop after `limit` rows. Vector search manages + // its own top_k limit in _configure_vector_search, so skip it here. + if (!_vector_search && lance_params.__isset.limit && lance_params.limit > 0) { + if (lance_scanner_set_limit(scanner, lance_params.limit) != 0) { + return _lance_error("set Lance scanner limit"); + } + } if (_vector_search) { RETURN_IF_ERROR(_configure_vector_search(scanner)); } diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScanNode.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScanNode.java index 76d222e7d59cb2..2e01e6462c11f5 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScanNode.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScanNode.java @@ -114,6 +114,19 @@ protected void doInitialize() throws UserException { } } + // A fragment-level LIMIT can be pushed into an ordinary Lance scan only when every predicate + // is already pushed into Lance (conjuncts is empty). Otherwise Doris re-filters the returned + // rows and truncating a fragment early could drop valid results. + // + // OFFSET needs no special handling: the Nereids SplitLimit rule rewrites Limit(limit, offset) + // into a global Limit(limit, offset) over a local Limit(limit + offset, 0), and the local + // bound is what lands on this scan node. So getLimit() already accounts for the offset and + // getOffset() is always 0 here; each fragment fetches up to limit + offset rows and the upper + // global LIMIT still applies the offset and the final bound. + private boolean canPushDownLimit() { + return hasLimit() && conjuncts.isEmpty(); + } + @Override protected void convertPredicate() { if (isExternalSearch()) { @@ -189,6 +202,12 @@ protected void setScanParams(TFileRangeDesc rangeDesc, Split split) { "Ordinary Lance scan split must contain one fragment"); } lanceParams.setFragmentIds(Collections.singletonList(lanceSplit.getFragmentId())); + // Push LIMIT into each fragment scanner only when it is safe to truncate a single + // fragment early. See canPushDownLimit(). Each scanner still returns at most `limit` + // rows and the upper LIMIT operator enforces the final bound across fragments. + if (canPushDownLimit()) { + lanceParams.setLimit(getLimit()); + } } TTableFormatFileDesc tableFormatParams = new TTableFormatFileDesc(); @@ -243,6 +262,9 @@ public String getNodeExplainString(String prefix, TExplainLevel detailLevel) { .append(((LanceExternalCatalog) lanceTable.getCatalog()).getLanceCatalogType()).append("\n"); result.append(prefix).append("lanceVersion=").append(plannedVersion).append("\n"); result.append(prefix).append("lanceFragments=").append(plannedFragments).append("\n"); + if (canPushDownLimit()) { + result.append(prefix).append("lanceLimit=").append(getLimit()).append("\n"); + } if (!lancePushdownPredicate.isEmpty()) { result.append(prefix).append("lancePushdownPredicate=") .append(lancePushdownPredicate).append("\n"); diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/LanceThriftContractTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/LanceThriftContractTest.java index e28dd3c73822ea..dc21fafa3148ab 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/LanceThriftContractTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/LanceThriftContractTest.java @@ -36,7 +36,8 @@ public void testLanceDescriptorCompactProtocolRoundTrip() throws Exception { TLanceFileDesc lanceDesc = new TLanceFileDesc() .setDatasetUri("s3://warehouse/db/table.lance") .setFragmentIds(Arrays.asList(7L, 11L)) - .setVersion(42L); + .setVersion(42L) + .setLimit(100L); TTableFormatFileDesc source = new TTableFormatFileDesc() .setTableFormatType(TableFormatType.LANCE.value()) .setLanceParams(lanceDesc); @@ -53,5 +54,27 @@ public void testLanceDescriptorCompactProtocolRoundTrip() throws Exception { Assert.assertEquals("s3://warehouse/db/table.lance", restored.getLanceParams().getDatasetUri()); Assert.assertEquals(Arrays.asList(7L, 11L), restored.getLanceParams().getFragmentIds()); Assert.assertEquals(42L, restored.getLanceParams().getVersion()); + Assert.assertTrue(restored.getLanceParams().isSetLimit()); + Assert.assertEquals(100L, restored.getLanceParams().getLimit()); + } + + @Test + public void testLanceDescriptorWithoutLimit() throws Exception { + TLanceFileDesc lanceDesc = new TLanceFileDesc() + .setDatasetUri("s3://warehouse/db/table.lance") + .setFragmentIds(Arrays.asList(1L)) + .setVersion(1L); + TTableFormatFileDesc source = new TTableFormatFileDesc() + .setTableFormatType(TableFormatType.LANCE.value()) + .setLanceParams(lanceDesc); + + TSerializer serializer = new TSerializer(new TCompactProtocol.Factory()); + byte[] bytes = serializer.serialize(source); + + TTableFormatFileDesc restored = new TTableFormatFileDesc(); + new TDeserializer(new TCompactProtocol.Factory()).deserialize(restored, bytes); + + // A scan without a pushable LIMIT must leave the field unset so the BE reads all rows. + Assert.assertFalse(restored.getLanceParams().isSetLimit()); } } diff --git a/gensrc/thrift/PlanNodes.thrift b/gensrc/thrift/PlanNodes.thrift index 3b7c377110ae08..d3f0583c2206a1 100644 --- a/gensrc/thrift/PlanNodes.thrift +++ b/gensrc/thrift/PlanNodes.thrift @@ -527,6 +527,10 @@ struct TLanceFileDesc { 1: optional string dataset_uri 2: optional list fragment_ids 3: optional i64 version + // Per-split row limit pushed down from the query LIMIT. Each scanner returns at + // most this many rows; the upper LIMIT operator still enforces the global bound. + // Only set for ordinary scans whose predicates are fully pushed into Lance. + 4: optional i64 limit } struct TTableFormatFileDesc {