Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion be/src/exec/operator/materialization_opertor.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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));
Expand Down
21 changes: 21 additions & 0 deletions be/src/exec/rowid_fetcher.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand Down Expand Up @@ -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();
}

Expand Down Expand Up @@ -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());
Expand All @@ -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++;
}

Expand All @@ -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));
Expand Down
2 changes: 2 additions & 0 deletions be/src/exec/rowid_fetcher.h
Original file line number Diff line number Diff line change
Expand Up @@ -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<std::string, TUnit::type> LanceFetchCountersProfile;
static const std::string TopNLazyMaterializationSecondPhaseLocalIOCount;
static const std::string TopNLazyMaterializationSecondPhaseLocalIOBytes;
static const std::string TopNLazyMaterializationSecondPhaseRemoteIOCount;
Expand Down Expand Up @@ -183,6 +184,7 @@ class RowIdStorageReader {
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;
};
Expand Down
13 changes: 13 additions & 0 deletions be/src/format_v2/table/lance_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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);
}
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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");
}
Expand Down Expand Up @@ -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");
Expand Down
49 changes: 49 additions & 0 deletions be/test/exec/operator/materialization_shared_state_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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 {

Expand All @@ -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");
Expand All @@ -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 {
Expand Down Expand Up @@ -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")));
Expand Down
60 changes: 60 additions & 0 deletions be/test/format_v2/table/lance_reader_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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());
}

Expand Down
36 changes: 36 additions & 0 deletions docs/lance-ann-profile.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P3] Correct the row-fetch count semantics for duplicate IDs

read_external_row_from_file_mapping deduplicates requested IDs before take_rows, then maps duplicate result positions back to the one fetched row. For a request containing the same ID twice, LanceRowIdFetchRows reports 1 while the result has 2 rows, contrary to this description. Describe the metric as distinct rows fetched per dataset/RPC, or count output copies separately.

| `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.
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
Loading
Loading