From 5b793d8801a8c049dd45037a2bcb86b2489b5d3a Mon Sep 17 00:00:00 2001 From: Gabriel Date: Thu, 1 Oct 2026 12:40:42 +0800 Subject: [PATCH 1/2] [improvement](lance) Expose query parallelism and complete row fetch profiles --- .../exec/operator/materialization_opertor.cpp | 6 +- be/src/exec/rowid_fetcher.cpp | 21 +++++++ be/src/exec/rowid_fetcher.h | 2 + be/src/format_v2/table/lance_reader.cpp | 13 ++++ .../materialization_shared_state_test.cpp | 37 ++++++++++++ be/test/format_v2/table/lance_reader_test.cpp | 60 +++++++++++++++++++ docs/lance-ann-profile.md | 36 +++++++++++ .../lance/source/LanceScanNode.java | 5 ++ .../VectorSearchTableValuedFunction.java | 11 +++- .../lance/source/LanceScanNodeTest.java | 6 +- .../VectorSearchTableValuedFunctionTest.java | 32 ++++++++++ gensrc/thrift/PlanNodes.thrift | 2 + .../lance/test_lance_vector_search.groovy | 20 +++++++ 13 files changed, 247 insertions(+), 4 deletions(-) diff --git a/be/src/exec/operator/materialization_opertor.cpp b/be/src/exec/operator/materialization_opertor.cpp index 47f65cd8c5fc00..16c2150ebf36f2 100644 --- a/be/src/exec/operator/materialization_opertor.cpp +++ b/be/src/exec/operator/materialization_opertor.cpp @@ -481,6 +481,9 @@ void MaterializationSharedState::_update_profile_info(int64_t backend_id, update_profile_info_key(RowIdStorageReader::LanceRowIdTakeReadTimeProfile, false); update_profile_info_key(RowIdStorageReader::LanceArrowToDorisBlockTimeProfile, false); update_profile_info_key(RowIdStorageReader::LanceRowIdFetchTotalTimeProfile, false); + for (const auto& [name, unit] : RowIdStorageReader::LanceFetchCountersProfile) { + update_profile_info_key(name, false); + } } Status MaterializationSharedState::create_muiltget_result(const Columns& columns, bool child_eos, @@ -674,7 +677,8 @@ Status MaterializationOperator::pull(RuntimeState* state, Block* output_block, b Status MaterializationOperator::push(RuntimeState* state, Block* in_block, bool eos) const { auto& local_state = get_local_state(state); - SCOPED_TIMER(local_state.exec_time_counter()); + // StatefulOperatorX::get_block_impl already times push() with this counter. + // Nesting the same timer counts the synchronous row-fetch RPC wait twice. if (!local_state._materialization_state.rpc_struct_inited) { RETURN_IF_ERROR(local_state._materialization_state.init_multi_requests( _materialization_node, state)); diff --git a/be/src/exec/rowid_fetcher.cpp b/be/src/exec/rowid_fetcher.cpp index 8c3b0023831899..92f8de5c2e6031 100644 --- a/be/src/exec/rowid_fetcher.cpp +++ b/be/src/exec/rowid_fetcher.cpp @@ -755,6 +755,12 @@ const std::string RowIdStorageReader::LanceRowIdTakeReadTimeProfile = "LanceRowI const std::string RowIdStorageReader::LanceArrowToDorisBlockTimeProfile = "LanceArrowToDorisBlockTime"; const std::string RowIdStorageReader::LanceRowIdFetchTotalTimeProfile = "LanceRowIdFetchTotalTime"; +const std::map RowIdStorageReader::LanceFetchCountersProfile = { + {"LanceRowIdFetchRows", TUnit::UNIT}, + {"LanceRowIdFetchCalls", TUnit::UNIT}, + {"LanceDataCacheBytesReadFromCache", TUnit::BYTES}, + {"LanceDataCacheBytesReadFromRemote", TUnit::BYTES}, +}; const std::string RowIdStorageReader::TopNLazyMaterializationSecondPhaseLocalIOCount = "TopNLazyMaterializationSecondPhaseLocalIOCount"; const std::string RowIdStorageReader::TopNLazyMaterializationSecondPhaseLocalIOBytes = @@ -889,6 +895,13 @@ Status RowIdStorageReader::read_lance_rows_by_row_ids( collect_lance_fetch_time(LanceRowIdTakeReadTimeProfile); collect_lance_fetch_time(LanceArrowToDorisBlockTimeProfile); collect_lance_fetch_time(LanceRowIdFetchTotalTimeProfile); + // close() publishes dataset-handle cache statistics; collect after it and preserve the + // units across the RPC instead of formatting byte/count values as nanoseconds. + for (const auto& [name, unit] : LanceFetchCountersProfile) { + if (const auto* counter = runtime_profile->get_counter(name); counter != nullptr) { + fetch_statistics->lance_fetch_counters.emplace(name, counter->value()); + } + } return Status::OK(); } @@ -1203,6 +1216,7 @@ Status RowIdStorageReader::read_batch_external_row( format_to(file_read_times_buffer, "["); std::map lance_fetch_times_ns; + std::map lance_fetch_counters; size_t idx = 0; for (const auto& [_, scan_info] : scan_rows) { format_to(file_read_lines_buffer, "{}, ", scan_info.first.size()); @@ -1213,6 +1227,9 @@ Status RowIdStorageReader::read_batch_external_row( for (const auto& [time_name, time_value] : fetch_statistics[idx].lance_fetch_times_ns) { lance_fetch_times_ns[time_name] += time_value; } + for (const auto& [name, value] : fetch_statistics[idx].lance_fetch_counters) { + lance_fetch_counters[name] += value; + } idx++; } @@ -1232,6 +1249,10 @@ Status RowIdStorageReader::read_batch_external_row( fmt::to_string(file_read_bytes_buffer)); runtime_profile->add_info_string(FileScannerV2::FileReadTimeProfile, fmt::to_string(file_read_times_buffer)); + for (const auto& [name, value] : lance_fetch_counters) { + runtime_profile->add_info_string( + name, PrettyPrinter::print(value, LanceFetchCountersProfile.at(name))); + } for (const auto& [time_name, time_value] : lance_fetch_times_ns) { runtime_profile->add_info_string(time_name, PrettyPrinter::print(time_value, TUnit::TIME_NS)); diff --git a/be/src/exec/rowid_fetcher.h b/be/src/exec/rowid_fetcher.h index ab8019e8f32641..fed8e5ccd5235b 100644 --- a/be/src/exec/rowid_fetcher.h +++ b/be/src/exec/rowid_fetcher.h @@ -108,6 +108,7 @@ class RowIdStorageReader { static const std::string LanceRowIdTakeReadTimeProfile; static const std::string LanceArrowToDorisBlockTimeProfile; static const std::string LanceRowIdFetchTotalTimeProfile; + static const std::map LanceFetchCountersProfile; static const std::string TopNLazyMaterializationSecondPhaseLocalIOCount; static const std::string TopNLazyMaterializationSecondPhaseLocalIOBytes; static const std::string TopNLazyMaterializationSecondPhaseRemoteIOCount; @@ -183,6 +184,7 @@ class RowIdStorageReader { int64_t init_reader_ms = 0; int64_t get_block_ms = 0; std::map lance_fetch_times_ns; + std::map lance_fetch_counters; std::string file_read_bytes; std::string file_read_times; }; diff --git a/be/src/format_v2/table/lance_reader.cpp b/be/src/format_v2/table/lance_reader.cpp index 52cf45533e5548..b9e5054a2a6e81 100644 --- a/be/src/format_v2/table/lance_reader.cpp +++ b/be/src/format_v2/table/lance_reader.cpp @@ -255,6 +255,10 @@ Status LanceTableReader::read_by_row_ids(const TFileRangeDesc& range, _row_id_fetch_total_time = ADD_CHILD_TIMER_WITH_LEVEL( _scanner_profile, "LanceRowIdFetchTotalTime", LANCE_READER_PROFILE, 1); } + auto* fetch_calls = ADD_CHILD_COUNTER_WITH_LEVEL(_scanner_profile, "LanceRowIdFetchCalls", + TUnit::UNIT, LANCE_READER_PROFILE, 1); + auto* fetch_rows = ADD_CHILD_COUNTER_WITH_LEVEL(_scanner_profile, "LanceRowIdFetchRows", + TUnit::UNIT, LANCE_READER_PROFILE, 1); SCOPED_TIMER(_row_id_fetch_total_time); // Phase-two row fetch does not execute FTS, so a reader created only for take_rows must not @@ -271,6 +275,7 @@ Status LanceTableReader::read_by_row_ids(const TFileRangeDesc& range, int32_t take_rows_status = 0; { SCOPED_TIMER(_row_id_take_read_time); + COUNTER_UPDATE(fetch_calls, 1); take_rows_status = lance_dataset_take_rows(_dataset, row_ids.data(), row_ids.size(), columns.data(), &stream); } @@ -313,6 +318,7 @@ Status LanceTableReader::read_by_row_ids(const TFileRangeDesc& range, record_batch, block, _global_rowid_context, &rows)); } fetched_rows += rows; + COUNTER_UPDATE(fetch_rows, rows); } if (fetched_rows != row_ids.size()) { return Status::InternalError("Lance row-id fetch returned {} rows for {} requested row ids", @@ -541,6 +547,9 @@ Status LanceTableReader::_validate_external_search_request() const { if (_search_kind == SearchKind::VECTOR && request.__isset.vector_search_options) { const auto& options = request.vector_search_options; + if (options.__isset.query_parallelism && options.query_parallelism < -1) { + return Status::InvalidArgument("Lance query_parallelism must be -1, 0, or positive"); + } if (options.__isset.nprobes && options.nprobes <= 0) { return Status::InvalidArgument("Lance nprobes must be positive"); } @@ -1076,6 +1085,10 @@ Status LanceTableReader::_configure_vector_search(LanceScanner* scanner, } if (request.__isset.vector_search_options) { const auto& options = request.vector_search_options; + if (options.__isset.query_parallelism && + lance_scanner_set_query_parallelism(scanner, options.query_parallelism) != 0) { + return lance_error("set Lance vector query parallelism"); + } if (options.__isset.nprobes && lance_scanner_set_nprobes(scanner, static_cast(options.nprobes)) != 0) { return lance_error("set Lance vector nprobes"); diff --git a/be/test/exec/operator/materialization_shared_state_test.cpp b/be/test/exec/operator/materialization_shared_state_test.cpp index 95ec24d0879fc0..b2a59817040d9c 100644 --- a/be/test/exec/operator/materialization_shared_state_test.cpp +++ b/be/test/exec/operator/materialization_shared_state_test.cpp @@ -27,6 +27,7 @@ #include "exec/operator/materialization_opertor.h" #include "exec/pipeline/dependency.h" #include "runtime/runtime_profile.h" +#include "testutil/mock/mock_runtime_state.h" namespace doris { @@ -43,6 +44,10 @@ void set_lance_fetch_profile(PMultiGetBlockV2* response_block, int64_t scale) { profile.add_info_string("LanceRowIdTakeReadTime", std::to_string(2 * scale) + "ns"); profile.add_info_string("LanceArrowToDorisBlockTime", std::to_string(3 * scale) + "ns"); profile.add_info_string("LanceRowIdFetchTotalTime", std::to_string(4 * scale) + "ns"); + profile.add_info_string("LanceRowIdFetchRows", std::to_string(5 * scale)); + profile.add_info_string("LanceRowIdFetchCalls", "1"); + profile.add_info_string("LanceDataCacheBytesReadFromCache", "64.00 KB"); + profile.add_info_string("LanceDataCacheBytesReadFromRemote", "32.00 KB"); profile.add_info_string("ScannersRunningTime", "0ms"); profile.add_info_string("InitReaderAvgTime", "0ms"); profile.add_info_string("GetBlockAvgTime", "0ms"); @@ -54,6 +59,34 @@ void set_lance_fetch_profile(PMultiGetBlockV2* response_block, int64_t scale) { } // namespace +TEST(MaterializationOperatorTimingTest, PushDoesNotDuplicateFrameworkExecTimer) { + ObjectPool pool; + MockRuntimeState state; + TPlanNode node; + node.__set_node_id(0); + node.__set_node_type(TPlanNodeType::MATERIALIZATION_NODE); + MaterializationOperator op(&pool, node, 0, state.desc_tbl()); + RuntimeProfile profile("materialization_timing"); + auto local = MaterializationLocalState::create_unique(&state, &op); + LocalStateInfo info {.parent_profile = &profile, + .scan_ranges = {}, + .shared_state = nullptr, + .shared_state_map = {}, + .task_idx = 0}; + ASSERT_TRUE(local->init(&state, info).ok()); + local->_materialization_state.rpc_struct_inited = true; + auto* timer = local->exec_time_counter(); + state.resize_op_id_to_local_state(-1); + state.emplace_local_state(0, std::move(local)); + Block block; + const auto before = timer->value(); + // The framework owns ExecTime. Direct push calls must not contribute a second sample. + for (int i = 0; i < 100; ++i) { + ASSERT_TRUE(op.push(&state, &block, false).ok()); + } + EXPECT_EQ(before, timer->value()); +} + class MaterializationSharedStateTest : public testing::Test { protected: void SetUp() override { @@ -316,6 +349,10 @@ TEST_F(MaterializationSharedStateTest, TestMergeMultiResponse) { EXPECT_EQ("20ns, ", fmt::to_string(backend1_info.at("LanceRowIdTakeReadTime"))); EXPECT_EQ("30ns, ", fmt::to_string(backend1_info.at("LanceArrowToDorisBlockTime"))); EXPECT_EQ("40ns, ", fmt::to_string(backend1_info.at("LanceRowIdFetchTotalTime"))); + EXPECT_EQ("50, ", fmt::to_string(backend1_info.at("LanceRowIdFetchRows"))); + EXPECT_EQ("1, ", fmt::to_string(backend1_info.at("LanceRowIdFetchCalls"))); + EXPECT_EQ("64.00 KB, ", fmt::to_string(backend1_info.at("LanceDataCacheBytesReadFromCache"))); + EXPECT_EQ("32.00 KB, ", fmt::to_string(backend1_info.at("LanceDataCacheBytesReadFromRemote"))); const auto& backend2_info = _shared_state->backend_profile_info_string.at(_backend_id2); EXPECT_EQ("1ns, ", fmt::to_string(backend2_info.at("LanceDatasetOpenTime"))); EXPECT_EQ("2ns, ", fmt::to_string(backend2_info.at("LanceRowIdTakeReadTime"))); diff --git a/be/test/format_v2/table/lance_reader_test.cpp b/be/test/format_v2/table/lance_reader_test.cpp index 93531688634ca2..e952d8da483c8f 100644 --- a/be/test/format_v2/table/lance_reader_test.cpp +++ b/be/test/format_v2/table/lance_reader_test.cpp @@ -1226,6 +1226,47 @@ TEST(LanceTableReaderVectorSearchTest, MultiVectorRejectsActualNullAndNonFiniteE } } +TEST(LanceTableReaderVectorSearchTest, RejectsInvalidQueryParallelismBeforeOpeningDataset) { + TQueryGlobals globals; + RuntimeState state(globals); + RuntimeProfile profile("lance_invalid_query_parallelism"); + auto params = make_float32_vector_search_params({0.0F, 0.0F, 0.0F}, 2, 0); + params.lance_scan_params.external_search_request.vector_search_options.__set_query_parallelism( + -2); + const Columns columns {projected_column("_distance", TYPE_FLOAT, true)}; + LanceTableReader reader; + const auto status = init_reader(&reader, columns, &state, &profile, ¶ms); + EXPECT_FALSE(status.ok()); + EXPECT_NE(std::string::npos, status.to_string().find("query_parallelism")); +} + +TEST(LanceTableReaderVectorSearchTest, QueryParallelismPreservesResults) { + const std::filesystem::path uri = "./be/test/format_v2/table/lance/data/all_types.lance"; + LanceFixtureInfo fixture; + ASSERT_TRUE(get_fixture_info(uri, &fixture).ok()); + const Columns columns {projected_column("row_id", TYPE_BIGINT, false), + projected_column("_distance", TYPE_FLOAT, true)}; + TQueryGlobals globals; + RuntimeState state(globals); + for (const int parallelism : {-1, 0, 1, 4}) { + SCOPED_TRACE(parallelism); + auto params = make_float32_vector_search_params({0.0F, 0.0F, 0.0F}, 2, 1); + params.lance_scan_params.external_search_request.vector_search_options + .__set_query_parallelism(parallelism); + RuntimeProfile profile("lance_query_parallelism"); + LanceTableReader reader; + ASSERT_TRUE(init_reader(&reader, columns, &state, &profile, ¶ms).ok()); + ASSERT_TRUE(prepare_fixture(&reader, uri, fixture, fixture.fragment_ids).ok()); + Block block; + add_output_columns(&block, columns); + const auto rows = read_vector_search_rows(&reader, &block); + ASSERT_EQ(2U, rows.size()); + EXPECT_EQ(2, rows[0].first); + EXPECT_EQ(4, rows[1].first); + ASSERT_TRUE(reader.close().ok()); + } +} + TEST(LanceTableReaderVectorSearchTest, SearchesWholeSnapshotWithOffsetAndDistance) { const std::filesystem::path dataset_uri = "./be/test/format_v2/table/lance/data/all_types.lance"; @@ -1431,6 +1472,25 @@ TEST(LanceTableReaderVectorSearchTest, ReturnsStableGlobalRowIdsAndFetchesPayloa EXPECT_EQ("extra", label_values.get_data_at(0).to_string()); EXPECT_EQ("unit-x", label_values.get_data_at(1).to_string()); EXPECT_EQ("extra", label_values.get_data_at(2).to_string()); + ASSERT_NE(nullptr, fetch_profile.get_counter("LanceRowIdFetchRows")); + EXPECT_EQ(3, fetch_profile.get_counter("LanceRowIdFetchRows")->value()); + ASSERT_NE(nullptr, fetch_profile.get_counter("LanceRowIdFetchCalls")); + EXPECT_EQ(1, fetch_profile.get_counter("LanceRowIdFetchCalls")->value()); + Block second_block; + add_output_columns(&second_block, payload_columns); + ASSERT_TRUE(payload_reader + .read_by_row_ids(make_lance_range(dataset_uri, fixture.version, + fixture.fragment_ids), + fetch_row_ids, &second_block) + .ok()); + EXPECT_EQ(6, fetch_profile.get_counter("LanceRowIdFetchRows")->value()); + EXPECT_EQ(2, fetch_profile.get_counter("LanceRowIdFetchCalls")->value()); + ASSERT_TRUE(payload_reader + .read_by_row_ids(make_lance_range(dataset_uri, fixture.version, + fixture.fragment_ids), + {}, &second_block) + .ok()); + EXPECT_EQ(2, fetch_profile.get_counter("LanceRowIdFetchCalls")->value()); EXPECT_TRUE(payload_reader.close().ok()); } diff --git a/docs/lance-ann-profile.md b/docs/lance-ann-profile.md index 51fdc15d095ef6..6a2a8bc92aa7f9 100644 --- a/docs/lance-ann-profile.md +++ b/docs/lance-ann-profile.md @@ -65,3 +65,39 @@ queries, inspect partition load together with execution bytes, requests, and partition cache misses. Use repeated queries and the operator-level elapsed times to assess latency; cumulative parallel stage times alone are not a critical path trace. + +## Query parallelism + +`vector_search` accepts an optional `"query_parallelism"` integer: + +- `0` (also the default when omitted): let Lance choose the parallelism. +- `-1`: use the available Lance CPU parallelism. +- Positive values: request that many concurrent partition searches, capped by + Lance's compute pool and execution-plan limits. + +For example, add `"query_parallelism" = "4"` alongside `"nprobes" = "64"`. +This controls concurrency inside a Lance search, independently of Doris scan +instances. Increasing it can increase intermediate candidates and memory usage; +measure both single-query latency and concurrent throughput. EXPLAIN displays an +explicit setting as `lanceQueryParallelism`. + +## Second-phase row-ID fetch + +Each `RowIDFetcher: BackendId:...` profile also reports: + +| Counter | Scope | +| --- | --- | +| `LanceRowIdFetchCalls` | Non-empty dataset `take_rows` calls. | +| `LanceRowIdFetchRows` | Rows converted from successful returned batches, including duplicate requested row IDs. | +| `LanceDataCacheBytesReadFromCache` | Logical data-file bytes served by Foyer for the dataset handles used by this fetch. | +| `LanceDataCacheBytesReadFromRemote` | Logical data-file bytes served through the Foyer origin path for those handles. | + +Byte counts are collected after closing the reader and summed across dataset +handles in the fetch RPC. They exclude block-alignment amplification, metadata, +and index reads. With Foyer disabled or for paths outside its data-file wrapper, +zero values do not imply zero physical IO. The current take API does not expose +physical request counts or a separate decode timer; scanner-plan IO counters +must not be substituted for them. + +`MATERIALIZATION_OPERATOR.ExecTime` includes its synchronous fetch wait once. +`MaxRpcTime` is nested within that execution time, not an additional duration. 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 a483ef52aa2d88..3f9505b87b9d9b 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 @@ -354,6 +354,11 @@ public String getNodeExplainString(String prefix, TExplainLevel detailLevel) { .append(scanPlan.vectorIndexStatus).append("\n"); result.append(prefix).append("lanceVectorColumn=") .append(vector.getColumn()).append("\n"); + if (externalSearchRequest.isSetVectorSearchOptions() + && externalSearchRequest.getVectorSearchOptions().isSetQueryParallelism()) { + result.append(prefix).append("lanceQueryParallelism=") + .append(externalSearchRequest.getVectorSearchOptions().getQueryParallelism()).append("\n"); + } result.append(prefix).append("lanceMetric=") .append(vector.isSetMetric() ? VectorSearchTableValuedFunction.metricName(vector.getMetric()) : "default") diff --git a/fe/fe-core/src/main/java/org/apache/doris/tablefunction/VectorSearchTableValuedFunction.java b/fe/fe-core/src/main/java/org/apache/doris/tablefunction/VectorSearchTableValuedFunction.java index c394adbe099947..93f7b2f4d560f4 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/tablefunction/VectorSearchTableValuedFunction.java +++ b/fe/fe-core/src/main/java/org/apache/doris/tablefunction/VectorSearchTableValuedFunction.java @@ -60,13 +60,14 @@ public class VectorSearchTableValuedFunction extends LanceExternalSearchTableVal private static final String ARROW_EXTENSION_NAME = "ARROW:extension:name"; private static final String QUERY_VECTOR = "query_vector"; private static final String METRIC = "metric"; + private static final String QUERY_PARALLELISM = "query_parallelism"; private static final String NPROBES = "nprobes"; private static final String REFINE_FACTOR = "refine_factor"; private static final String EF = "ef"; private static final String USE_INDEX = "use_index"; private static final Set PROPERTIES = ImmutableSet.of( TABLE, COLUMN, QUERY_VECTOR, TOP_K, OFFSET, METRIC, FILTER, - NPROBES, REFINE_FACTOR, EF, USE_INDEX); + NPROBES, REFINE_FACTOR, EF, USE_INDEX, QUERY_PARALLELISM); public VectorSearchTableValuedFunction(Map properties) throws AnalysisException { @@ -125,10 +126,16 @@ private static PreparedSearch prepare(Map properties, boolean de common, vectorFieldId, searchRequest, DISTANCE_COLUMN, "vector search"); } - private static TVectorSearchOptions buildVectorSearchOptions( + @VisibleForTesting + static TVectorSearchOptions buildVectorSearchOptions( Map params, boolean useIndex) throws AnalysisException { TVectorSearchOptions options = new TVectorSearchOptions(); boolean configured = false; + if (params.containsKey(QUERY_PARALLELISM)) { + options.setQueryParallelism((int) parseLong( + params.get(QUERY_PARALLELISM), QUERY_PARALLELISM, -1, Integer.MAX_VALUE)); + configured = true; + } if (params.containsKey(NPROBES)) { options.setNprobes(parsePositiveInt(params.get(NPROBES), NPROBES)); configured = true; 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..b68ca1f437cfde 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 @@ -563,10 +563,14 @@ public void testExternalSearchUsesOneSplitPerIndexSegmentAndKeepsUnindexedFragme new LanceIndexSegmentInfo(secondSegment, "vector_idx", Collections.singletonList(9), Arrays.asList(3L, 4L), IndexType.VECTOR, "L2"))); - LanceScanNode node = newSearchNode(metadata, vectorSearchRequest(5, 0)); + TExternalSearchRequest request = vectorSearchRequest(5, 0); + request.setVectorSearchOptions(new TVectorSearchOptions().setQueryParallelism(4)); + LanceScanNode node = newSearchNode(metadata, request); List splits = node.getSplits(3); + Assert.assertTrue(node.getNodeExplainString("", TExplainLevel.NORMAL) + .contains("lanceQueryParallelism=4")); Assert.assertTrue(node.getNodeExplainString("", TExplainLevel.NORMAL) .contains("lanceVectorIndexStatus=USED")); diff --git a/fe/fe-core/src/test/java/org/apache/doris/tablefunction/VectorSearchTableValuedFunctionTest.java b/fe/fe-core/src/test/java/org/apache/doris/tablefunction/VectorSearchTableValuedFunctionTest.java index 1fe7caefa3355e..c8c3da2c6c0aa8 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/tablefunction/VectorSearchTableValuedFunctionTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/tablefunction/VectorSearchTableValuedFunctionTest.java @@ -22,18 +22,50 @@ import org.apache.doris.common.AnalysisException; import org.apache.doris.datasource.lance.metadata.LanceTableAccess; import org.apache.doris.datasource.lance.metadata.LanceTableMetadata; +import org.apache.doris.thrift.TVectorSearchOptions; import org.apache.arrow.vector.types.pojo.ArrowType; import org.apache.arrow.vector.types.pojo.Field; import org.apache.arrow.vector.types.pojo.Schema; +import org.apache.thrift.TDeserializer; +import org.apache.thrift.TSerializer; import org.junit.Assert; import org.junit.Test; import java.nio.charset.StandardCharsets; import java.util.Arrays; import java.util.Collections; +import java.util.Map; public class VectorSearchTableValuedFunctionTest { + private TVectorSearchOptions parseOptions(Map params) throws Exception { + return VectorSearchTableValuedFunction.buildVectorSearchOptions(params, true); + } + + @Test + public void testQueryParallelismOptions() throws Exception { + Assert.assertNull(parseOptions(Collections.emptyMap())); + Assert.assertFalse(parseOptions(Collections.singletonMap("nprobes", "4")).isSetQueryParallelism()); + for (String value : new String[] {"-1", "0", "1", "4", "2147483647"}) { + TVectorSearchOptions options = parseOptions(Collections.singletonMap("query_parallelism", value)); + Assert.assertNotNull("query_parallelism must be serialized", options); + TVectorSearchOptions decoded = new TVectorSearchOptions(); + new TDeserializer().deserialize(decoded, new TSerializer().serialize(options)); + Assert.assertTrue(decoded.isSetQueryParallelism()); + Assert.assertEquals(Integer.parseInt(value), decoded.getQueryParallelism()); + Assert.assertFalse(decoded.isSetNprobes()); + } + } + + @Test + public void testRejectInvalidQueryParallelism() { + for (String value : new String[] {"-2", "2147483648", "1.5", "abc", ""}) { + AnalysisException error = Assert.assertThrows(AnalysisException.class, + () -> parseOptions(Collections.singletonMap("query_parallelism", value))); + Assert.assertTrue(error.getMessage(), error.getMessage().contains("query_parallelism")); + } + } + @Test public void testParseQuotedMultiLevelNamespace() throws AnalysisException { TableName tableName = VectorSearchTableValuedFunction.parseTableName( diff --git a/gensrc/thrift/PlanNodes.thrift b/gensrc/thrift/PlanNodes.thrift index 3e626e6f0ad3eb..6e4fa8fd77b1af 100644 --- a/gensrc/thrift/PlanNodes.thrift +++ b/gensrc/thrift/PlanNodes.thrift @@ -552,6 +552,8 @@ struct TVectorSearchOptions { 2: optional i32 refine_factor 3: optional i32 ef 4: optional bool use_index + // Lance: -1 uses available CPU parallelism, 0 selects automatically, positive values cap it. + 5: optional i32 query_parallelism } // The active union field identifies the logical search kind. A future hybrid field can contain both diff --git a/regression-test/suites/external_table_p0/lance/test_lance_vector_search.groovy b/regression-test/suites/external_table_p0/lance/test_lance_vector_search.groovy index 7215db90265131..d90acafef7bf7e 100644 --- a/regression-test/suites/external_table_p0/lance/test_lance_vector_search.groovy +++ b/regression-test/suites/external_table_p0/lance/test_lance_vector_search.groovy @@ -148,6 +148,26 @@ suite("test_lance_vector_search", "p0,external") { ORDER BY _distance, row_id """ + def baselineRows = sql "SELECT row_id, label, _distance FROM ${indexedTopFive} ORDER BY _distance, row_id" + [-1, 0, 1, 2, 4].each { parallelism -> + String parallelSearch = indexedTopFive.substring(0, indexedTopFive.length() - 1) + + ", \"query_parallelism\"=\"${parallelism}\")" + explain { + sql("SELECT row_id, label, _distance FROM ${parallelSearch}") + contains "lanceQueryParallelism=${parallelism}" + } + assertEquals(baselineRows, + sql("SELECT row_id, label, _distance FROM ${parallelSearch} ORDER BY _distance, row_id")) + } + ["-2", "2147483648", "1.5", "invalid"].each { invalid -> + String invalidSearch = indexedTopFive.substring(0, indexedTopFive.length() - 1) + + ", \"query_parallelism\"=\"${invalid}\")" + test { + sql "SELECT row_id FROM ${invalidSearch}" + exception "query_parallelism" + } + } + // Disable the index explicitly for the exact flat baseline. qt_flat_l2_topk """ SELECT row_id, label, _distance From bfbbbaabae51fab3f76f3b191c1f862b3f598647 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Thu, 1 Oct 2026 19:03:37 +0800 Subject: [PATCH 2/2] [fix](test) Initialize materialization timing test row descriptors --- .../operator/materialization_shared_state_test.cpp | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/be/test/exec/operator/materialization_shared_state_test.cpp b/be/test/exec/operator/materialization_shared_state_test.cpp index b2a59817040d9c..b59ba2540ad020 100644 --- a/be/test/exec/operator/materialization_shared_state_test.cpp +++ b/be/test/exec/operator/materialization_shared_state_test.cpp @@ -62,9 +62,21 @@ void set_lance_fetch_profile(PMultiGetBlockV2* response_block, int64_t scale) { TEST(MaterializationOperatorTimingTest, PushDoesNotDuplicateFrameworkExecTimer) { ObjectPool pool; MockRuntimeState state; + // Even an empty-block timing test needs a real tuple: OperatorXBase builds + // a RowDescriptor from the plan and requires non-empty, resolvable tuple IDs. + TTupleDescriptor tuple; + tuple.__set_id(0); + TDescriptorTable thrift_desc; + thrift_desc.__set_tupleDescriptors({tuple}); + DescriptorTbl* desc_tbl = nullptr; + ASSERT_TRUE(DescriptorTbl::create(&pool, thrift_desc, &desc_tbl).ok()); + state.set_desc_tbl(desc_tbl); + TPlanNode node; node.__set_node_id(0); node.__set_node_type(TPlanNodeType::MATERIALIZATION_NODE); + node.__set_row_tuples({0}); + node.__set_nullable_tuples({false}); MaterializationOperator op(&pool, node, 0, state.desc_tbl()); RuntimeProfile profile("materialization_timing"); auto local = MaterializationLocalState::create_unique(&state, &op);