This is an automated email from the ASF dual-hosted git repository.

yiguolei pushed a commit to branch branch-4.2
in repository https://gitbox.apache.org/repos/asf/doris.git

commit 5320dddf5dc6d96b476f111f4484be83c30828a5
Author: Gabriel <[email protected]>
AuthorDate: Thu Oct 8 09:25:47 2026 +0800

    [improvement](lance) Expose query parallelism and complete row fetch 
profiles (branch-4.1) (#68690)
    
    ### What problem does this PR solve?
    
    Lance vector searches cannot configure the SDK's query parallelism from
    Doris, and second-phase row-ID fetch profiles omit available read
    counters. Materialization also nests two timers for the same `ExecTime`
    counter, counting the synchronous fetch wait twice.
    
    This PR targets `branch-4.1` and:
    
    - Adds optional `vector_search("query_parallelism" = "4", ...)` support
    through FE validation, Thrift, and the existing Lance-C API. Accepts
    `-1`, `0`, and positive integers; omitting it preserves the SDK default.
    EXPLAIN shows an explicit setting.
    - Carries row-ID fetch call counts, returned row counts, and per-dataset
    Foyer cache/origin logical read bytes through the fetch RPC into each
    `RowIDFetcher` profile, preserving their units. Statistics are captured
    after the reader closes.
    - Removes the nested materialization execution timer; the operator
    framework continues to count the fetch wait once.
    
    The read-byte counters exclude metadata/index IO and block-alignment
    amplification. The current take API does not expose physical request
    counts or a separate decode timer; zero Foyer counters on bypassed paths
    do not mean zero physical IO. These boundaries are documented in
    `docs/lance-ann-profile.md`.
    
    No third-party dependency or patch changes are required.
    
    ### Validation
    
    - 43 focused FE tests passed, including option bounds, omitted defaults,
    Thrift round trips, and EXPLAIN propagation.
    - FE Checkstyle passed with zero violations.
    - Clang-format 16 checks passed for every changed C++ source/header.
    - C++ syntax checks passed for all changed translation units and both
    affected BE test files using regenerated Thrift definitions.
    - Groovy compilation passed for the changed regression suite.
    - Added BE tests for parameter validation, result preservation,
    row-fetch counters, RPC profile propagation, and avoiding duplicate
    execution timing.
    - Added external regression coverage comparing indexed results across
    parallelism settings and rejecting invalid values. Full BE UT and
    external regression execution remain for CI; local C++ checks were
    syntax-only.
    
    ### Release note
    
    Add configurable Lance vector query parallelism and row-ID fetch read
    counters, and correct materialization execution-time accounting.
    
    ### Check List (For Author)
    
    - Test
      - [x] Regression test
      - [x] Unit Test
    - Behavior changed:
    - [x] Yes. Explicit query parallelism is now accepted; existing defaults
    are preserved. Materialization profiles no longer double-count fetch
    execution time.
    - Does this need documentation?
      - [x] Yes. Updated `docs/lance-ann-profile.md` in this PR.
    
    ### Check List (For Reviewer who merge this PR)
    
    - [ ] Confirm the release note
    - [ ] Confirm test cases
    - [ ] Confirm document
    - [ ] Add branch pick label
---
 be/src/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 +++++
 .../operator/materialization_shared_state_test.cpp | 49 ++++++++++++++++++
 be/test/format_v2/table/lance_reader_test.cpp      | 60 ++++++++++++++++++++++
 docs/lance-ann-profile.md                          | 36 +++++++++++++
 .../datasource/lance/source/LanceScanNode.java     |  5 ++
 .../VectorSearchTableValuedFunction.java           | 11 +++-
 .../datasource/lance/source/LanceScanNodeTest.java |  6 ++-
 .../VectorSearchTableValuedFunctionTest.java       | 32 ++++++++++++
 gensrc/thrift/PlanNodes.thrift                     |  2 +
 .../lance/test_lance_vector_search.groovy          | 20 ++++++++
 13 files changed, 259 insertions(+), 4 deletions(-)

diff --git a/be/src/exec/operator/materialization_opertor.cpp 
b/be/src/exec/operator/materialization_opertor.cpp
index 47f65cd8c5f..16c2150ebf3 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 8c3b0023831..92f8de5c2e6 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<std::string, TUnit::type> 
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<std::string, int64_t> lance_fetch_times_ns;
+        std::map<std::string, int64_t> 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 ab8019e8f32..fed8e5ccd52 100644
--- a/be/src/exec/rowid_fetcher.h
+++ b/be/src/exec/rowid_fetcher.h
@@ -108,6 +108,7 @@ public:
     static const std::string LanceRowIdTakeReadTimeProfile;
     static const std::string LanceArrowToDorisBlockTimeProfile;
     static const std::string LanceRowIdFetchTotalTimeProfile;
+    static const std::map<std::string, TUnit::type> LanceFetchCountersProfile;
     static const std::string TopNLazyMaterializationSecondPhaseLocalIOCount;
     static const std::string TopNLazyMaterializationSecondPhaseLocalIOBytes;
     static const std::string TopNLazyMaterializationSecondPhaseRemoteIOCount;
@@ -183,6 +184,7 @@ private:
         int64_t init_reader_ms = 0;
         int64_t get_block_ms = 0;
         std::map<std::string, int64_t> lance_fetch_times_ns;
+        std::map<std::string, int64_t> 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 52cf45533e5..b9e5054a2a6 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<uint32_t>(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 95ec24d0879..b59ba2540ad 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,46 @@ void set_lance_fetch_profile(PMultiGetBlockV2* 
response_block, int64_t scale) {
 
 } // namespace
 
+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);
+    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 +361,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 93531688634..e952d8da483 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, 
&params);
+    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, 
&params).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 51fdc15d095..6a2a8bc92aa 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 a483ef52aa2..3f9505b87b9 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 class LanceScanNode extends FileQueryScanNode {
                         .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 c394adbe099..93f7b2f4d56 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<String> 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<String, String> properties)
             throws AnalysisException {
@@ -125,10 +126,16 @@ public class VectorSearchTableValuedFunction extends 
LanceExternalSearchTableVal
                 common, vectorFieldId, searchRequest, DISTANCE_COLUMN, "vector 
search");
     }
 
-    private static TVectorSearchOptions buildVectorSearchOptions(
+    @VisibleForTesting
+    static TVectorSearchOptions buildVectorSearchOptions(
             Map<String, String> 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 6131a2fcb2c..b68ca1f437c 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 class LanceScanNodeTest {
                         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<Split> 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 1fe7caefa33..c8c3da2c6c0 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.catalog.Column;
 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<String, String> 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 3e626e6f0ad..6e4fa8fd77b 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 7215db90265..d90acafef7b 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


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to