diff --git a/configs/milvus.yaml b/configs/milvus.yaml index 56c52012f25..fd6524a1002 100644 --- a/configs/milvus.yaml +++ b/configs/milvus.yaml @@ -621,7 +621,8 @@ queryNode: scalarIndex: disable vectorField: disable # cache warmup for vector field raw data is by default disabled. vectorIndex: disable - lazyManifestReaderEnabled: false # Defer the Storage V3 projected ChunkReader/translator/cache tree for warmup=disable fields without load-time side effects. + lazyManifestReaderEnabled: false # Defer Storage V3 manifest readers for warmup=disable fields until first access. + lazyJsonStatsEnabled: true # Defer eligible Storage V3 JSON stats metadata until first access. # If evictionEnabled is true, a background thread will run every evictionIntervalMs to determine if an # eviction is necessary and the amount of data to evict from memory/disk. # - If the current memory/disk usage exceeds the high watermark, an eviction will be triggered to evict data from memory/disk diff --git a/docs/design-docs/design_docs/segcore/20260722-storage-v3-manifest-lazy-materialization.md b/docs/design-docs/design_docs/segcore/20260722-storage-v3-manifest-lazy-materialization.md index cce8bcdca62..c811a1cba5f 100644 --- a/docs/design-docs/design_docs/segcore/20260722-storage-v3-manifest-lazy-materialization.md +++ b/docs/design-docs/design_docs/segcore/20260722-storage-v3-manifest-lazy-materialization.md @@ -1,4 +1,4 @@ -# Storage V3 Manifest Task 延迟物化设计 +# Storage V3 Manifest Task 与 JSON Stats 延迟物化设计 ## 1. 目标 @@ -12,13 +12,17 @@ - `ChunkedColumnGroup`; - 真实 `ProxyChunkColumn`。 +独立的 `lazyJsonStatsEnabled` 开关允许将 JSON key stats 顶层对象的远端 `meta.json`、 +Parquet metadata、projected reader 和 cache tree 初始化延迟到首个可使用 stats 的表达式。 +Ready 阶段只发布绑定当前 generation 的 `LazyJsonStats` facade。 + 设计边界如下: - 保持 `SegmentLoadInfo` 已有 Task、projection 和 eager/lazy fallback 粒度; - 不因开启本功能而按字段重新拆分 Column Group; - 支持普通字段、RowID、Timestamp、INT64 PK 和 VARCHAR PK; - Ready 阶段保持 manifest Task 延迟,查询期需要数据时允许正常物化; -- 正确性优先,不增加 JSON、Tantivy、PK 或 MVCC 专用的“保冷”执行路径; +- 正确性优先,不增加 Tantivy、PK 或 MVCC 专用的“保冷”执行路径; - 保持 Storage V1/V2、external collection 和配置关闭时的行为不变。 ## 2. 生效条件 @@ -36,6 +40,11 @@ Task 仅在以下条件同时满足时进入延迟物化: Task,避免同一 generation 内出现部分 eager、部分 lazy 的混合状态。已发布 Task 不随 动态配置变化而改变。 +JSON stats 不依赖 `lazyManifestReaderEnabled`。它在 +`queryNode.segcore.tieredStorage.lazyJsonStatsEnabled=true`、上述第 2 至第 4 个条件成立, +且 effective scalar-field warmup policy 为 `disable` 时延迟物化。`sync`、`async`、 +JSON stats 配置关闭以及非 Storage V3 segment 继续在 Load 阶段初始化,保持既有行为。 + ## 3. Task 不变量 Lazy 开关不参与 Task 分组。每个 Task 保持 `SegmentLoadInfo` 已确定的字段集合和 @@ -55,15 +64,21 @@ Ready 阶段的对象关系如下: ```text RuntimeResourceState ├── generation runtime Reader -└── fields - └── LazyManifestProxyColumn - └── Task-scoped LazyManifestColumnGroup - └── ManifestColumnGroupBuildContext - ├── generation Reader 强引用 - ├── 原始 column-group index - ├── Task projection - ├── FieldMeta 快照 - └── mmap、priority、cache key 等 Translator 输入 +├── fields +│ └── LazyManifestProxyColumn +│ └── Task-scoped LazyManifestColumnGroup +│ └── ManifestColumnGroupBuildContext +│ ├── generation Reader 强引用 +│ ├── 原始 column-group index +│ ├── Task projection +│ ├── FieldMeta 快照 +│ └── mmap、priority、cache key 等 Translator 输入 +└── json_stats + └── LazyJsonStats + └── JsonStatsBuildContext + ├── generation FileManagerContext + ├── files、base path、mmap 与 priority + └── generation insert channel 与 resolved warmup policy ``` 每个 Load/Reopen generation 创建并发布自己的 runtime Reader。Lazy Task 直接强引用 @@ -101,6 +116,9 @@ Materialize(op_ctx) - 失败时不发布半成品 group; - 成功后所有 facade 复用同一 group。 +`LazyJsonStats` 使用相同的 single-flight、独立取消和失败重试语义。成功前不发布 +`JsonKeyStats`;原有 `internal_json_stats_latency_load` 只记录真正发生的物化耗时。 + ## 6. Facade 接口语义 `LazyManifestProxyColumn` 保存共享 Task、FieldId、FieldMeta 和行数。 @@ -125,6 +143,9 @@ Materialize(op_ctx) `CellsLoaded()` 仅表示对应 cache cell 是否已经加载。真实 group 已创建但 cell 尚未加载 时仍可返回 false,因此它可以作为非物化状态探针。 +JSON stats facade 的存在性检查只在表达式已满足开关、非空 path 且 path 不含数组下标后 +触发物化。首个符合条件的表达式加载 stats,之后同 generation 的查询复用同一实例。 + ## 7. Lazy 资格 Task 内任一字段具有以下 Load 阶段副作用时,整个 Task 使用原 eager 路径: @@ -187,6 +208,7 @@ Column Group identity 包含: - 旧 snapshot 可以继续读取旧 generation; - 新 facade 不捕获 committer、可变 runtime 或当前 PublishedState; +- JSON stats facade 捕获本 generation 的 FileManagerContext、insert channel 与配置; - Reopen 失败不发布 staged state; - 同一 generation 内全部新 Task 使用同一个配置快照。 @@ -210,19 +232,20 @@ queryNode: segcore: tieredStorage: lazyManifestReaderEnabled: false + lazyJsonStatsEnabled: true ``` -配置默认关闭并支持动态刷新。Go paramtable 通过 C bridge 更新 `SegcoreConfig` 中的原子 -布尔值。 +Manifest reader 延迟加载默认关闭,JSON stats 延迟加载默认开启;两个开关相互独立且均 +支持动态刷新。Go paramtable 通过 C bridge 更新 `SegcoreConfig` 中各自的原子布尔值。 ## 12. 兼容性 -- 配置关闭时走原 eager 路径; +- 对应配置关闭时各自走原 eager 路径; - Storage V1/V2 和无 ManifestPath segment 不进入本路径; - external collection 保持既有加载方式; - 不修改 milvus-storage API; - 不改变既有 Task/projection 划分; -- 不为 JSON、Tantivy 或其他索引增加专用的“不物化”执行分支; +- JSON stats 仅延迟顶层初始化,查询执行逻辑与 raw-data fallback 选择不变; - 查询链路需要 chunk 布局、大小或原始数据时正常物化。 ## 13. 验证 @@ -233,6 +256,8 @@ queryNode: - Ready 阶段 facade 与未缓存的 PK/Timestamp slot; - 多字段 Task 的 sibling 共享和并发 single-flight; - cancellation 与非取消失败后的重试; +- JSON stats 的 Ready 冷态、并发首访 single-flight、取消和失败重试; +- JSON stats 的 `disable`、`sync`、配置关闭与非 Storage V3 边界; - `DataByteSize()`、布局接口和真实读取的按需物化; - RowID、INT64/VARCHAR PK、Timestamp 与 commit timestamp; - Reopen generation 重绑和 schema-only drop; diff --git a/internal/core/src/exec/expression/Expr.h b/internal/core/src/exec/expression/Expr.h index 9d83350cd35..a4c4b78f616 100644 --- a/internal/core/src/exec/expression/Expr.h +++ b/internal/core/src/exec/expression/Expr.h @@ -2745,12 +2745,12 @@ class SegmentExpr : public Expr { return false; } - // Check whether this expression can use JsonStats without pinning. - // All conditions are available before execution path determination. + // Check cheap expression constraints before the presence check, which may + // materialize lazy JsonStats on the first eligible query. bool CanUseJsonStatsAtInit() const { - return plan_options_.expr_use_json_stats && HasJsonStats(field_id_) && - !nested_path_.empty() && !PathContainsInteger(nested_path_); + return plan_options_.expr_use_json_stats && !nested_path_.empty() && + !PathContainsInteger(nested_path_) && HasJsonStats(field_id_); } virtual bool diff --git a/internal/core/src/segcore/ChunkedSegmentSealedImpl.cpp b/internal/core/src/segcore/ChunkedSegmentSealedImpl.cpp index 11e90f010ab..55881824aaf 100644 --- a/internal/core/src/segcore/ChunkedSegmentSealedImpl.cpp +++ b/internal/core/src/segcore/ChunkedSegmentSealedImpl.cpp @@ -168,8 +168,207 @@ struct ChunkedSegmentSealedImpl::ColumnSizeEstimateState { storagev2translator::ColumnSizeEstimateResult estimate_; }; +struct LazyJsonStats::State { + struct BuildAttempt { + std::condition_variable cv; + bool done{false}; + std::exception_ptr error; + }; + + State(int64_t segment_id, FieldId field_id, Loader loader) + : segment_id(segment_id), + field_id(field_id), + loader(std::move(loader)) { + } + + State(int64_t segment_id, + FieldId field_id, + std::shared_ptr stats) + : segment_id(segment_id), field_id(field_id), stats(std::move(stats)) { + } + + int64_t segment_id; + FieldId field_id; + Loader loader; + mutable std::mutex mutex; + std::shared_ptr stats; + std::shared_ptr attempt; +}; + +LazyJsonStats::LazyJsonStats(int64_t segment_id, + FieldId field_id, + Loader loader) + : state_(std::make_unique(segment_id, field_id, std::move(loader))) { +} + +LazyJsonStats::LazyJsonStats(int64_t segment_id, + FieldId field_id, + std::shared_ptr stats) + : state_(std::make_unique(segment_id, field_id, std::move(stats))) { +} + +LazyJsonStats::~LazyJsonStats() = default; + +std::shared_ptr +LazyJsonStats::Materialize(milvus::OpContext* op_ctx) { + std::shared_ptr attempt; + Loader loader; + while (true) { + CheckCancellation(op_ctx, + state_->segment_id, + state_->field_id.get(), + "LazyJsonStats::Materialize()"); + std::unique_lock lock(state_->mutex); + if (state_->stats != nullptr) { + return state_->stats; + } + AssertInfo(state_->loader != nullptr, + "json stats loader is null, segment {}, field {}", + state_->segment_id, + state_->field_id.get()); + if (state_->attempt != nullptr) { + attempt = state_->attempt; + while (!attempt->done) { + attempt->cv.wait_for(lock, std::chrono::milliseconds(20)); + CheckCancellation(op_ctx, + state_->segment_id, + state_->field_id.get(), + "LazyJsonStats::Materialize()"); + } + if (attempt->error != nullptr) { + std::rethrow_exception(attempt->error); + } + if (state_->stats != nullptr) { + return state_->stats; + } + continue; + } + + attempt = std::make_shared(); + state_->attempt = attempt; + loader = state_->loader; + break; + } + + std::shared_ptr stats; + std::exception_ptr error; + bool cancelled = false; + try { + stats = loader(op_ctx); + AssertInfo(stats != nullptr, + "json stats loader returned null, segment {}, field {}", + state_->segment_id, + state_->field_id.get()); + CheckCancellation(op_ctx, + state_->segment_id, + state_->field_id.get(), + "LazyJsonStats::Materialize()"); + } catch (const SegcoreError& e) { + cancelled = e.get_error_code() == ErrorCode::FollyCancel; + error = std::current_exception(); + } catch (...) { + error = std::current_exception(); + } + + { + std::lock_guard lock(state_->mutex); + if (error == nullptr) { + state_->stats = stats; + state_->loader = nullptr; + } else if (!cancelled) { + attempt->error = error; + } + attempt->done = true; + if (state_->attempt == attempt) { + state_->attempt.reset(); + } + } + attempt->cv.notify_all(); + if (error != nullptr) { + std::rethrow_exception(error); + } + return stats; +} + +std::shared_ptr +LazyJsonStats::MaterializedIfReady() const { + std::lock_guard lock(state_->mutex); + return state_->stats; +} + namespace { +struct JsonStatsBuildContext { + int64_t segment_id; + FieldId field_id; + int64_t build_id; + int64_t version; + int file_count; + std::string base_path; + bool enable_mmap; + int64_t stats_size; + storage::FileManagerContext file_manager_context; + milvus::Config config; +}; + +std::shared_ptr +CreateJsonKeyStats(const JsonStatsBuildContext& context, + milvus::OpContext* op_ctx) { + CheckCancellation(op_ctx, + context.segment_id, + context.field_id.get(), + "LazyJsonStats::Materialize()"); + LOG_INFO( + "start load json key stats, segment:{}, field:{}, build:{}, " + "version:{}, file_count:{}, base_path:{}, enable_mmap:{}, " + "stats_size:{}", + context.segment_id, + context.field_id.get(), + context.build_id, + context.version, + context.file_count, + context.base_path, + context.enable_mmap, + context.stats_size); + + auto index = std::make_shared( + context.file_manager_context, true); + milvus::tracer::TraceContext trace_ctx; + try { + milvus::ScopedTimer timer( + "json_stats_load", + [](double us) { + milvus::monitor::internal_json_stats_latency_load.Observe( + us / 1000.0); + }, + milvus::ScopedTimer::LogLevel::Info); + index->Load(trace_ctx, context.config); + CheckCancellation(op_ctx, + context.segment_id, + context.field_id.get(), + "LazyJsonStats::Materialize()"); + } catch (const std::exception& e) { + LOG_WARN( + "failed load json key stats, segment:{}, field:{}, build:{}, " + "version:{}, error:{}", + context.segment_id, + context.field_id.get(), + context.build_id, + context.version, + e.what()); + throw; + } + + LOG_INFO( + "load json key stats success, segment:{}, field:{}, build:{}, " + "version:{}", + context.segment_id, + context.field_id.get(), + context.build_id, + context.version); + return index; +} + struct ManifestColumnGroupBuildContext { int64_t segment_id; int64_t original_column_group_index; @@ -2405,10 +2604,10 @@ ChunkedSegmentSealedImpl::GetJsonStats(milvus::OpContext* op_ctx, return nullptr; } auto iter = runtime->json_stats.find(field_id); - if (iter == runtime->json_stats.end()) { + if (iter == runtime->json_stats.end() || iter->second == nullptr) { return nullptr; } - return iter->second; + return iter->second->Materialize(op_ctx); } void @@ -6483,11 +6682,14 @@ ChunkedSegmentSealedImpl::BuildTextIndexFromFiles( return cache_slot; } -std::shared_ptr +std::shared_ptr ChunkedSegmentSealedImpl::BuildJsonKeyStatsIndex( milvus::OpContext* op_ctx, const std::shared_ptr& - info_proto) { + info_proto, + const SegmentLoadInfo& segment_load_info, + const SchemaPtr& schema_snapshot, + bool lazy_json_stats_enabled) { auto field_id = milvus::FieldId(info_proto->fieldid()); CheckCancellation(op_ctx, id_, @@ -6505,19 +6707,6 @@ ChunkedSegmentSealedImpl::BuildJsonKeyStatsIndex( return nullptr; } - LOG_INFO( - "start load json key stats, segment:{}, field:{}, build:{}, " - "version:{}, " - "file_count:{}, base_path:{}, enable_mmap:{}, stats_size:{}", - id_, - info_proto->fieldid(), - info_proto->buildid(), - info_proto->version(), - info_proto->files_size(), - info_proto->base_path(), - info_proto->enable_mmap(), - info_proto->stats_size()); - milvus::storage::FieldDataMeta field_data_meta{info_proto->collectionid(), info_proto->partitionid(), this->get_segment_id(), @@ -6552,42 +6741,57 @@ ChunkedSegmentSealedImpl::BuildJsonKeyStatsIndex( if (!info_proto->base_path().empty()) { config[STATS_BASE_PATH_KEY] = info_proto->base_path(); } - auto load_info_snapshot = CaptureLoadInfoSnapshot(); - config[JSON_STATS_CACHE_SHARD_KEY] = load_info_snapshot->GetInsertChannel(); + config[JSON_STATS_CACHE_SHARD_KEY] = segment_load_info.GetInsertChannel(); milvus::storage::FileManagerContext file_ctx( field_data_meta, index_meta, remote_chunk_manager, fs); - auto index = std::make_shared(file_ctx, true); - milvus::tracer::TraceContext trace_ctx; - try { - milvus::ScopedTimer timer( - "json_stats_load", - [](double us) { - milvus::monitor::internal_json_stats_latency_load.Observe( - us / 1000.0); - }, - milvus::ScopedTimer::LogLevel::Info); - index->Load(trace_ctx, config); - } catch (std::exception& e) { - LOG_WARN( - "failed load json key stats, segment:{}, field:{}, build:{}, " - "version:{}, error:{}", + auto effective_warmup = getCacheWarmupPolicy(info_proto->warmup_policy(), + /*is_vector=*/false, + /*is_index=*/false, + /*in_load_list=*/true); + const bool can_defer = + lazy_json_stats_enabled && + segment_load_info.GetStorageVersion() == STORAGE_V3 && + segment_load_info.HasManifestPath() && + !schema_snapshot->is_external_collection() && + effective_warmup == CacheWarmupPolicy::CacheWarmupPolicy_Disable; + if (can_defer) { + // Bind the resolved policy to this generation. A later global config + // refresh must not make a published cold facade warm its child caches. + config[milvus::index::WARMUP] = "disable"; + } + + JsonStatsBuildContext build_context{ + .segment_id = id_, + .field_id = field_id, + .build_id = info_proto->buildid(), + .version = info_proto->version(), + .file_count = info_proto->files_size(), + .base_path = info_proto->base_path(), + .enable_mmap = info_proto->enable_mmap(), + .stats_size = info_proto->stats_size(), + .file_manager_context = std::move(file_ctx), + .config = std::move(config), + }; + auto loader = [context = std::move(build_context)]( + milvus::OpContext* materialize_ctx) { + return CreateJsonKeyStats(context, materialize_ctx); + }; + auto stats = + std::make_shared(id_, field_id, std::move(loader)); + if (can_defer) { + LOG_INFO( + "defer load json key stats until first access, segment:{}, " + "field:{}, build:{}, version:{}", id_, info_proto->fieldid(), info_proto->buildid(), - info_proto->version(), - e.what()); - throw; + info_proto->version()); + return stats; } - LOG_INFO( - "load json key stats success, segment:{}, field:{}, build:{}, " - "version:{}", - id_, - info_proto->fieldid(), - info_proto->buildid(), - info_proto->version()); - return index; + stats->Materialize(op_ctx); + return stats; } void @@ -6597,19 +6801,25 @@ ChunkedSegmentSealedImpl::LoadBatchJsonKeyIndexes( FieldId, std::shared_ptr>& infos, const SchemaPtr& schema_snapshot, + const SegmentLoadInfo& segment_load_info, + bool lazy_json_stats_enabled, StagedStateCommitter& committer) { for (const auto& [field_id, info_proto] : infos) { AssertInfo(field_exists_in_schema(schema_snapshot, field_id), "field {} not found in schema when loading json stats", field_id.get()); - auto index = BuildJsonKeyStatsIndex(op_ctx, info_proto); - if (index == nullptr) { + auto stats = BuildJsonKeyStatsIndex(op_ctx, + info_proto, + segment_load_info, + schema_snapshot, + lazy_json_stats_enabled); + if (stats == nullptr) { continue; } committer.Commit( - [field_id = field_id, index = std::move(index)]( + [field_id = field_id, stats = std::move(stats)]( RuntimeResourceState& runtime, PublishedSegmentState&) mutable { - runtime.json_stats[field_id] = std::move(index); + runtime.json_stats[field_id] = std::move(stats); }); } } @@ -8135,6 +8345,10 @@ ChunkedSegmentSealedImpl::PrepareLoadDiffForReopen( const SchemaPtr& schema_snapshot, StagedStateCommitter& committer) { milvus::tracer::TraceContext trace_ctx; + const bool lazy_manifest_reader_enabled = + segcore_config_.get_lazy_manifest_reader_enabled(); + const bool lazy_json_stats_enabled = + segcore_config_.get_lazy_json_stats_enabled(); CheckCancellation(op_ctx, id_, "ChunkedSegmentSealedImpl::ApplyLoadDiff()"); if (!diff.indexes_to_load.empty()) { @@ -8178,8 +8392,6 @@ ChunkedSegmentSealedImpl::PrepareLoadDiffForReopen( *milvus::storage::LoonFFIPropertiesSingleton::GetInstance() .GetProperties()); auto column_groups = segment_load_info.GetColumnGroups(); - const bool lazy_manifest_reader_enabled = - segcore_config_.get_lazy_manifest_reader_enabled(); auto arrow_schema = schema_snapshot->ConvertToLoonArrowSchema( /*text_lob_as_binary=*/true); auto needed_columns = std::make_shared>(); @@ -8285,12 +8497,20 @@ ChunkedSegmentSealedImpl::PrepareLoadDiffForReopen( CheckCancellation(op_ctx, id_, "ChunkedSegmentSealedImpl::ApplyLoadDiff()"); if (!diff.json_stats_to_load.empty()) { - LoadBatchJsonKeyIndexes( - op_ctx, diff.json_stats_to_load, schema_snapshot, committer); + LoadBatchJsonKeyIndexes(op_ctx, + diff.json_stats_to_load, + schema_snapshot, + segment_load_info, + lazy_json_stats_enabled, + committer); } if (!diff.json_stats_to_replace.empty()) { - LoadBatchJsonKeyIndexes( - op_ctx, diff.json_stats_to_replace, schema_snapshot, committer); + LoadBatchJsonKeyIndexes(op_ctx, + diff.json_stats_to_replace, + schema_snapshot, + segment_load_info, + lazy_json_stats_enabled, + committer); } CheckCancellation(op_ctx, id_, "ChunkedSegmentSealedImpl::ApplyLoadDiff()"); diff --git a/internal/core/src/segcore/ChunkedSegmentSealedImpl.h b/internal/core/src/segcore/ChunkedSegmentSealedImpl.h index c9bfce78cee..f6082bfd81c 100644 --- a/internal/core/src/segcore/ChunkedSegmentSealedImpl.h +++ b/internal/core/src/segcore/ChunkedSegmentSealedImpl.h @@ -99,6 +99,33 @@ class PkIndexCell; class TimestampData; class TimestampIndex; +// Generation-bound facade for JSON key stats. The facade itself is published +// in RuntimeResourceState; the potentially expensive remote metadata and +// parquet reader initialization happens on first eligible query access. +class LazyJsonStats { + public: + using Loader = + std::function(milvus::OpContext*)>; + + LazyJsonStats(int64_t segment_id, FieldId field_id, Loader loader); + + LazyJsonStats(int64_t segment_id, + FieldId field_id, + std::shared_ptr stats); + + ~LazyJsonStats(); + + std::shared_ptr + Materialize(milvus::OpContext* op_ctx); + + std::shared_ptr + MaterializedIfReady() const; + + private: + struct State; + std::unique_ptr state_; +}; + using namespace milvus::cachinglayer; // Test-only accessor that simulates v2/v3 segment state (raw timestamp column @@ -352,8 +379,7 @@ class ChunkedSegmentSealedImpl : public SegmentSealed { std::unordered_map text_lob_paths; std::unordered_map text_indexes; std::vector json_indices; - std::unordered_map> - json_stats; + std::unordered_map> json_stats; std::shared_ptr reader; std::shared_ptr timestamps; std::shared_ptr timestamp_index; @@ -1699,11 +1725,14 @@ class ChunkedSegmentSealedImpl : public SegmentSealed { milvus::OpContext* op_ctx = nullptr, PublishMode publish_mode = PublishMode::Drain); - std::shared_ptr + std::shared_ptr BuildJsonKeyStatsIndex( milvus::OpContext* op_ctx, const std::shared_ptr& - info_proto); + info_proto, + const SegmentLoadInfo& segment_load_info, + const SchemaPtr& schema_snapshot, + bool lazy_json_stats_enabled); void LoadBatchJsonKeyIndexes( @@ -1713,6 +1742,8 @@ class ChunkedSegmentSealedImpl : public SegmentSealed { std::shared_ptr>& infos, const SchemaPtr& schema_snapshot, + const SegmentLoadInfo& segment_load_info, + bool lazy_json_stats_enabled, StagedStateCommitter& committer); template @@ -2690,7 +2721,58 @@ class ChunkedSegmentSealedImpl : public SegmentSealed { auto current = CapturePublishedState(); auto next = ClonePublishedState(current); auto runtime = CloneRuntimeResourceState(current->runtime); - runtime->json_stats[field_id] = std::move(stats); + runtime->json_stats[field_id] = + std::make_shared(id_, field_id, std::move(stats)); + next->runtime = ToConstRuntimeState(std::move(runtime)); + NormalizePublishedState(*next); + PublishStateOnline(std::move(next)); + } + + void + SetLazyJsonStatsForTesting(FieldId field_id, LazyJsonStats::Loader loader) { + std::lock_guard reopen_guard(reopen_mutex_); + auto current = CapturePublishedState(); + auto next = ClonePublishedState(current); + auto runtime = CloneRuntimeResourceState(current->runtime); + runtime->json_stats[field_id] = + std::make_shared(id_, field_id, std::move(loader)); + next->runtime = ToConstRuntimeState(std::move(runtime)); + NormalizePublishedState(*next); + PublishStateOnline(std::move(next)); + } + + bool + JsonStatsMaterializedForTesting(FieldId field_id) const { + auto runtime = CaptureRuntimeResourceState(); + auto it = runtime->json_stats.find(field_id); + return it != runtime->json_stats.end() && it->second != nullptr && + it->second->MaterializedIfReady() != nullptr; + } + + void + LoadJsonStatsForTesting( + const std::shared_ptr& + info, + const SegmentLoadInfo& segment_load_info, + bool lazy_json_stats_enabled) { + std::lock_guard reopen_guard(reopen_mutex_); + auto current = CapturePublishedState(); + auto next = ClonePublishedState(current); + auto runtime = CloneRuntimeResourceState(current->runtime); + next->load_info = + std::make_shared(segment_load_info); + next->runtime = ToConstRuntimeState(runtime); + StagedStateCommitter committer(*this, runtime.get(), next.get()); + std::unordered_map< + FieldId, + std::shared_ptr> + infos{{FieldId(info->fieldid()), info}}; + LoadBatchJsonKeyIndexes(nullptr, + infos, + current->schema, + segment_load_info, + lazy_json_stats_enabled, + committer); next->runtime = ToConstRuntimeState(std::move(runtime)); NormalizePublishedState(*next); PublishStateOnline(std::move(next)); diff --git a/internal/core/src/segcore/SegcoreConfig.h b/internal/core/src/segcore/SegcoreConfig.h index 05f6c0f6ada..ae31d2d924c 100644 --- a/internal/core/src/segcore/SegcoreConfig.h +++ b/internal/core/src/segcore/SegcoreConfig.h @@ -230,6 +230,16 @@ class SegcoreConfig { return lazy_manifest_reader_enabled_.load(std::memory_order_relaxed); } + void + set_lazy_json_stats_enabled(bool value) { + lazy_json_stats_enabled_.store(value, std::memory_order_relaxed); + } + + bool + get_lazy_json_stats_enabled() const { + return lazy_json_stats_enabled_.load(std::memory_order_relaxed); + } + void set_reject_remote_vector_output(bool value) { reject_remote_vector_output_ = value; @@ -312,6 +322,7 @@ class SegcoreConfig { inline static bool enable_gis_split_fusion_ = false; inline static bool prefer_field_data_when_index_has_raw_data_ = false; inline static std::atomic lazy_manifest_reader_enabled_ = false; + inline static std::atomic lazy_json_stats_enabled_ = true; inline static bool reject_remote_vector_output_ = false; inline static std::atomic take_for_output_result_count_limit_{ kDefaultTakeForOutputResultCountLimit}; diff --git a/internal/core/src/segcore/segcore_init_c.cpp b/internal/core/src/segcore/segcore_init_c.cpp index a503a3547c0..52f45ac8957 100644 --- a/internal/core/src/segcore/segcore_init_c.cpp +++ b/internal/core/src/segcore/segcore_init_c.cpp @@ -140,6 +140,13 @@ SegcoreSetLazyManifestReaderEnabled(const bool value) { config.set_lazy_manifest_reader_enabled(value); } +extern "C" void +SegcoreSetLazyJsonStatsEnabled(const bool value) { + milvus::segcore::SegcoreConfig& config = + milvus::segcore::SegcoreConfig::default_config(); + config.set_lazy_json_stats_enabled(value); +} + extern "C" void SegcoreSetNlist(const int64_t value) { milvus::segcore::SegcoreConfig& config = diff --git a/internal/core/src/segcore/segcore_init_c.h b/internal/core/src/segcore/segcore_init_c.h index 8f0209268c7..6dd22fdd474 100644 --- a/internal/core/src/segcore/segcore_init_c.h +++ b/internal/core/src/segcore/segcore_init_c.h @@ -118,6 +118,9 @@ SegcoreGetTakeForOutputResultCountLimit(); void SegcoreSetLazyManifestReaderEnabled(const bool value); +void +SegcoreSetLazyJsonStatsEnabled(const bool value); + void SegcoreCloseGlog(); diff --git a/internal/core/unittest/test_sealed.cpp b/internal/core/unittest/test_sealed.cpp index 42b1dfef1af..45133416385 100644 --- a/internal/core/unittest/test_sealed.cpp +++ b/internal/core/unittest/test_sealed.cpp @@ -81,6 +81,7 @@ #include "segcore/SegmentLoadInfo.h" #include "segcore/SegmentSealed.h" #include "segcore/Types.h" +#include "segcore/default_fs.h" #include "segcore/storagev2translator/SystemIndexTranslator.h" #include "storage/FileManager.h" #include "storage/InsertData.h" @@ -167,6 +168,47 @@ CreateWarmupPolicySchema(bool include_vector) { return schema; } +std::shared_ptr +MakeUnloadedJsonStats(FieldId field_id) { + proto::schema::FieldSchema field_schema; + field_schema.set_fieldid(field_id.get()); + field_schema.set_name("payload"); + field_schema.set_data_type(proto::schema::DataType::JSON); + storage::FieldDataMeta field_data_meta{/*collection_id=*/1, + /*partition_id=*/2, + /*segment_id=*/3, + field_id.get(), + field_schema}; + storage::IndexMeta index_meta{ + /*segment_id=*/3, field_id.get(), /*build_id=*/4, /*version=*/5}; + auto chunk_manager = storage::RemoteChunkManagerSingleton::GetInstance() + .GetRemoteChunkManager(); + storage::FileManagerContext file_context( + field_data_meta, + index_meta, + chunk_manager, + milvus::segcore::GetDefaultArrowFileSystem()); + return std::make_shared(file_context, true); +} + +std::shared_ptr +MakeInvalidJsonStatsLoadInfo(FieldId field_id, + const FieldMeta& field_meta, + const std::string& warmup_policy) { + auto info = std::make_shared(); + info->set_collectionid(1); + info->set_partitionid(2); + info->set_fieldid(field_id.get()); + info->set_buildid(4); + info->set_version(5); + info->set_stats_size(64); + info->set_base_path("/unused-json-stats-test-path"); + info->set_warmup_policy(warmup_policy); + info->add_files("unsupported.stats"); + *info->mutable_schema() = field_meta.ToProto(); + return info; +} + std::shared_ptr MakeWarmupTestColumnGroups() { auto column_groups = std::make_shared(); @@ -5732,6 +5774,200 @@ TEST(SealedSegmentCowState, JsonIndexReplaceNgramWithScalarErasesNgramPath) { original_index); } +TEST(SealedSegmentCowState, LazyJsonStatsConcurrentFirstAccessLoadsOnce) { + auto schema = std::make_shared(); + auto pk = schema->AddDebugField("pk", DataType::INT64); + auto json = schema->AddDebugField("payload", DataType::JSON); + schema->set_primary_field_id(pk); + + auto segment = CreateSealedSegment(schema, nullptr, 7001); + auto* sealed = dynamic_cast(segment.get()); + ASSERT_NE(sealed, nullptr); + auto stats = MakeUnloadedJsonStats(json); + + std::atomic load_calls{0}; + sealed->SetLazyJsonStatsForTesting( + json, [stats, &load_calls](milvus::OpContext*) { + load_calls.fetch_add(1, std::memory_order_relaxed); + std::this_thread::sleep_for(std::chrono::milliseconds(30)); + return stats; + }); + EXPECT_FALSE(sealed->JsonStatsMaterializedForTesting(json)); + + constexpr int kThreadCount = 16; + std::atomic ready{0}; + std::atomic start{false}; + std::atomic failed{false}; + std::vector workers; + workers.reserve(kThreadCount); + for (int i = 0; i < kThreadCount; ++i) { + workers.emplace_back([&]() { + ready.fetch_add(1, std::memory_order_acq_rel); + while (!start.load(std::memory_order_acquire)) { + std::this_thread::yield(); + } + try { + if (sealed->GetJsonStats(nullptr, json) != stats) { + failed.store(true, std::memory_order_release); + } + } catch (...) { + failed.store(true, std::memory_order_release); + } + }); + } + while (ready.load(std::memory_order_acquire) != kThreadCount) { + std::this_thread::yield(); + } + start.store(true, std::memory_order_release); + for (auto& worker : workers) { + worker.join(); + } + + EXPECT_FALSE(failed.load(std::memory_order_acquire)); + EXPECT_EQ(load_calls.load(std::memory_order_relaxed), 1); + EXPECT_TRUE(sealed->JsonStatsMaterializedForTesting(json)); + EXPECT_EQ(sealed->GetJsonStats(nullptr, json), stats); + EXPECT_EQ(load_calls.load(std::memory_order_relaxed), 1); +} + +TEST(SealedSegmentCowState, LazyJsonStatsCancellationAllowsFreshRetry) { + auto schema = std::make_shared(); + auto pk = schema->AddDebugField("pk", DataType::INT64); + auto json = schema->AddDebugField("payload", DataType::JSON); + schema->set_primary_field_id(pk); + + auto segment = CreateSealedSegment(schema, nullptr, 7002); + auto* sealed = dynamic_cast(segment.get()); + ASSERT_NE(sealed, nullptr); + auto stats = MakeUnloadedJsonStats(json); + + std::atomic load_calls{0}; + sealed->SetLazyJsonStatsForTesting( + json, [stats, &load_calls](milvus::OpContext*) { + load_calls.fetch_add(1, std::memory_order_relaxed); + return stats; + }); + + folly::CancellationSource source; + source.requestCancellation(); + milvus::OpContext cancelled_ctx(source.getToken()); + try { + (void)sealed->GetJsonStats(&cancelled_ctx, json); + FAIL() << "expected cancelled JSON stats materialization"; + } catch (const SegcoreError& error) { + EXPECT_EQ(error.get_error_code(), ErrorCode::FollyCancel); + } + EXPECT_EQ(load_calls.load(std::memory_order_relaxed), 0); + EXPECT_FALSE(sealed->JsonStatsMaterializedForTesting(json)); + + EXPECT_EQ(sealed->GetJsonStats(nullptr, json), stats); + EXPECT_EQ(load_calls.load(std::memory_order_relaxed), 1); + EXPECT_TRUE(sealed->JsonStatsMaterializedForTesting(json)); +} + +TEST(SealedSegmentCowState, LazyJsonStatsFailureIsRetryable) { + auto schema = std::make_shared(); + auto pk = schema->AddDebugField("pk", DataType::INT64); + auto json = schema->AddDebugField("payload", DataType::JSON); + schema->set_primary_field_id(pk); + + auto segment = CreateSealedSegment(schema, nullptr, 7003); + auto* sealed = dynamic_cast(segment.get()); + ASSERT_NE(sealed, nullptr); + auto stats = MakeUnloadedJsonStats(json); + + std::atomic load_calls{0}; + sealed->SetLazyJsonStatsForTesting( + json, [stats, &load_calls](milvus::OpContext*) { + if (load_calls.fetch_add(1, std::memory_order_relaxed) == 0) { + throw std::runtime_error("injected JSON stats load failure"); + } + return stats; + }); + + EXPECT_THROW((void)sealed->GetJsonStats(nullptr, json), std::runtime_error); + EXPECT_FALSE(sealed->JsonStatsMaterializedForTesting(json)); + EXPECT_EQ(sealed->GetJsonStats(nullptr, json), stats); + EXPECT_EQ(load_calls.load(std::memory_order_relaxed), 2); + EXPECT_TRUE(sealed->JsonStatsMaterializedForTesting(json)); +} + +TEST(SealedSegmentCowState, LazyJsonStatsHonorsWarmupAndFeatureGate) { + auto schema = std::make_shared(); + auto pk = schema->AddDebugField("pk", DataType::INT64); + auto json = schema->AddDebugField("payload", DataType::JSON); + schema->set_primary_field_id(pk); + + proto::segcore::SegmentLoadInfo load_proto; + load_proto.set_segmentid(7004); + load_proto.set_partitionid(2); + load_proto.set_collectionid(1); + load_proto.set_num_of_rows(1); + load_proto.set_storageversion(STORAGE_V3); + load_proto.set_manifest_path("unused-manifest-path"); + load_proto.set_insert_channel("lazy-json-stats-test-channel"); + SegmentLoadInfo segment_load_info(load_proto, schema); + + auto lazy_segment = CreateSealedSegment(schema, nullptr, 7004); + auto* lazy = dynamic_cast(lazy_segment.get()); + ASSERT_NE(lazy, nullptr); + auto disabled_warmup = + MakeInvalidJsonStatsLoadInfo(json, schema->operator[](json), "disable"); + EXPECT_NO_THROW(lazy->LoadJsonStatsForTesting( + disabled_warmup, segment_load_info, /*lazy_enabled=*/true)); + EXPECT_FALSE(lazy->JsonStatsMaterializedForTesting(json)); + EXPECT_THROW((void)lazy->GetJsonStats(nullptr, json), SegcoreError); + EXPECT_FALSE(lazy->JsonStatsMaterializedForTesting(json)); + + auto sync_segment = CreateSealedSegment(schema, nullptr, 7004); + auto* sync = dynamic_cast(sync_segment.get()); + ASSERT_NE(sync, nullptr); + auto sync_warmup = + MakeInvalidJsonStatsLoadInfo(json, schema->operator[](json), "sync"); + EXPECT_THROW(sync->LoadJsonStatsForTesting(sync_warmup, + segment_load_info, + /*lazy_enabled=*/true), + SegcoreError); + EXPECT_FALSE(sync->JsonStatsMaterializedForTesting(json)); + + auto disabled_segment = CreateSealedSegment(schema, nullptr, 7004); + auto* disabled = + dynamic_cast(disabled_segment.get()); + ASSERT_NE(disabled, nullptr); + EXPECT_THROW(disabled->LoadJsonStatsForTesting(disabled_warmup, + segment_load_info, + /*lazy_enabled=*/false), + SegcoreError); + EXPECT_FALSE(disabled->JsonStatsMaterializedForTesting(json)); + + proto::segcore::SegmentLoadInfo storage_v2_proto = load_proto; + storage_v2_proto.set_storageversion(STORAGE_V2); + SegmentLoadInfo storage_v2_load_info(storage_v2_proto, schema); + auto storage_v2_segment = CreateSealedSegment(schema, nullptr, 7004); + auto* storage_v2 = + dynamic_cast(storage_v2_segment.get()); + ASSERT_NE(storage_v2, nullptr); + EXPECT_THROW(storage_v2->LoadJsonStatsForTesting(disabled_warmup, + storage_v2_load_info, + /*lazy_enabled=*/true), + SegcoreError); + EXPECT_FALSE(storage_v2->JsonStatsMaterializedForTesting(json)); + + auto external_schema = std::make_shared(*schema); + external_schema->set_external_source("s3://unused-json-stats-test"); + external_schema->set_external_spec(R"({"format":"parquet"})"); + SegmentLoadInfo external_load_info(load_proto, external_schema); + auto external_segment = CreateSealedSegment(external_schema, nullptr, 7004); + auto* external = + dynamic_cast(external_segment.get()); + ASSERT_NE(external, nullptr); + EXPECT_THROW(external->LoadJsonStatsForTesting(disabled_warmup, + external_load_info, + /*lazy_enabled=*/true), + SegcoreError); + EXPECT_FALSE(external->JsonStatsMaterializedForTesting(json)); +} + TEST(SealedSegmentCowState, JsonStatsLivesInRuntimeSnapshot) { auto schema = std::make_shared(); auto pk = schema->AddDebugField("pk", DataType::INT64); diff --git a/internal/streamingcoord/server/balancer/channel/manager.go b/internal/streamingcoord/server/balancer/channel/manager.go index 80480406105..187b9e17ae1 100644 --- a/internal/streamingcoord/server/balancer/channel/manager.go +++ b/internal/streamingcoord/server/balancer/channel/manager.go @@ -42,7 +42,7 @@ type ( StreamingVersion *streamingpb.StreamingVersion Version typeutil.VersionInt64Pair CChannelAssignment *streamingpb.CChannelAssignment - PChannelView *PChannelView + PChannels []string Relations []types.PChannelInfoAssigned ShardAssignments map[int64]types.ShardAssignmentInfo ReplicateConfiguration *commonpb.ReplicateConfiguration @@ -738,14 +738,15 @@ func (cm *ChannelManager) getNewIncomingTask(newConfig *replicateutil.ConfigHelp func (cm *ChannelManager) applyAssignments(cb WatchChannelAssignmentsCallback) (typeutil.VersionInt64Pair, error) { cm.cond.L.Lock() assignments := make([]types.PChannelInfoAssigned, 0, len(cm.channels)) + pchannels := make([]string, 0, len(cm.channels)) for _, c := range cm.channels { + pchannels = append(pchannels, c.Name()) if c.IsAssigned() { assignments = append(assignments, c.CurrentAssignment()) } } version := cm.version cchannelAssignment := proto.Clone(cm.cchannelMeta).(*streamingpb.CChannelMeta) - pchannelViews := newPChannelView(cm.channels) shardAssignmentProvider := cm.shardAssignmentProvider cm.cond.L.Unlock() @@ -760,7 +761,7 @@ func (cm *ChannelManager) applyAssignments(cb WatchChannelAssignmentsCallback) ( CChannelAssignment: &streamingpb.CChannelAssignment{ Meta: cchannelAssignment, }, - PChannelView: pchannelViews, + PChannels: pchannels, Relations: assignments, ShardAssignments: shardAssignments, ReplicateConfiguration: replicateConfig, diff --git a/internal/streamingcoord/server/balancer/channel/manager_test.go b/internal/streamingcoord/server/balancer/channel/manager_test.go index aa9fa5ff543..55fdbb2c896 100644 --- a/internal/streamingcoord/server/balancer/channel/manager_test.go +++ b/internal/streamingcoord/server/balancer/channel/manager_test.go @@ -159,6 +159,7 @@ func TestChannelManager(t *testing.T) { param, err := m.GetLatestChannelAssignment() oldLocalVersion := param.Version.Local assert.NoError(t, err) + assert.ElementsMatch(t, []string{"test-channel"}, param.PChannels) assert.Equal(t, m.ReplicateRole(), replicateutil.RolePrimary) // Test update replicate configurations @@ -692,6 +693,7 @@ func TestChannelManagerWatch(t *testing.T) { go func() { defer close(done) err := manager.WatchAssignmentResult(ctx, func(param WatchChannelAssignmentsCallbackParam) error { + assert.ElementsMatch(t, []string{"test-channel"}, param.PChannels) select { case called <- struct{}{}: default: diff --git a/internal/streamingcoord/server/service/assignment.go b/internal/streamingcoord/server/service/assignment.go index cc7ed06daff..ad05111e6fd 100644 --- a/internal/streamingcoord/server/service/assignment.go +++ b/internal/streamingcoord/server/service/assignment.go @@ -324,10 +324,8 @@ func (s *assignmentServiceImpl) buildForcePromoteConfiguration(ctx context.Conte return nil, nil, status.NewInvalidArgument("force promote requires current cluster in existing configuration; cluster %s not found in config", currentClusterID) } - // Get pchannels from PChannelView for validation - pchannels := lo.MapToSlice(latestAssignment.PChannelView.Channels, func(_ channel.ChannelID, ch *channel.PChannelMeta) string { - return ch.Name() - }) + // Copy before sorting because the assignment snapshot may be shared with other readers. + pchannels := append([]string(nil), latestAssignment.PChannels...) // Sort pchannels for consistent ordering (map iteration order is randomized) sort.Strings(pchannels) diff --git a/internal/streamingcoord/server/service/assignment_test.go b/internal/streamingcoord/server/service/assignment_test.go index 82d7b7e4528..d9f49a8b9e6 100644 --- a/internal/streamingcoord/server/service/assignment_test.go +++ b/internal/streamingcoord/server/service/assignment_test.go @@ -88,11 +88,7 @@ func TestAssignmentService(t *testing.T) { // Test illegal replicate configuration cfg := &commonpb.ReplicateConfiguration{} b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, }, nil).Maybe() _, err = as.UpdateReplicateConfiguration(context.Background(), &streamingpb.UpdateReplicateConfigurationRequest{ Configuration: cfg, @@ -119,11 +115,7 @@ func TestAssignmentService(t *testing.T) { // Test idempotent b.EXPECT().GetLatestChannelAssignment().Unset() b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, ReplicateConfiguration: cfg, }, nil).Maybe() _, err = as.UpdateReplicateConfiguration(context.Background(), &streamingpb.UpdateReplicateConfigurationRequest{ @@ -145,11 +137,7 @@ func TestAssignmentService(t *testing.T) { // Test update on secondary path, it should be block until the replicate configuration is changed. b.EXPECT().GetLatestChannelAssignment().Unset() b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, ReplicateConfiguration: &commonpb.ReplicateConfiguration{ Clusters: []*commonpb.MilvusCluster{ {ClusterId: "by-dev", Pchannels: []string{"by-dev-1"}, ConnectionParam: &commonpb.ConnectionParam{Uri: "http://test:19530", Token: "by-dev"}}, @@ -267,11 +255,7 @@ func TestForcePromoteOnPrimaryCluster(t *testing.T) { b.EXPECT().Close().Return().Maybe() // Current cluster is a primary b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, }, nil).Maybe() balance.Register(b) @@ -320,11 +304,7 @@ func TestForcePromoteSuccess(t *testing.T) { }, } b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, ReplicateConfiguration: currentReplicateConfig, }, nil).Maybe() balance.Register(b) @@ -391,11 +371,7 @@ func TestForcePromoteIdempotent(t *testing.T) { }, } b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, ReplicateConfiguration: alreadyPromotedConfig, }, nil).Maybe() balance.Register(b) @@ -546,11 +522,7 @@ func TestForcePromoteBroadcastOtherError(t *testing.T) { b.EXPECT().Close().Return().Maybe() b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, }, nil).Maybe() balance.Register(b) @@ -597,11 +569,7 @@ func TestForcePromoteBroadcastAppendError(t *testing.T) { }, } b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, ReplicateConfiguration: currentReplicateConfig, }, nil).Maybe() balance.Register(b) @@ -701,11 +669,7 @@ func TestUpdateReplicateConfigNonPrimaryBroadcastError(t *testing.T) { }).Maybe() b.EXPECT().Close().Return().Maybe() b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, }, nil).Maybe() balance.Register(b) @@ -757,11 +721,7 @@ func TestUpdateReplicateConfigBroadcastError(t *testing.T) { }).Maybe() b.EXPECT().Close().Return().Maybe() b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, }, nil).Maybe() balance.Register(b) @@ -834,20 +794,12 @@ func TestUpdateReplicateConfigSecondValidateSameConfig(t *testing.T) { if callCount <= 1 { // First call: config is different (nil) return &balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, }, nil } // Second call: config is now the same (was applied by another path) return &balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, ReplicateConfiguration: cfg, }, nil }) @@ -1055,11 +1007,7 @@ func TestHandleForcePromoteSameConfigAfterBroadcasterCheck(t *testing.T) { // Return the standalone primary config directly // This tests the idempotent path where config already matches b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, ReplicateConfiguration: forcePromoteCfg, }, nil).Maybe() balance.Register(b) @@ -1106,15 +1054,10 @@ func TestHandleForcePromoteValidatorError(t *testing.T) { return ctx.Err() }).Maybe() b.EXPECT().Close().Return().Maybe() - // Return config with extra pchannel in PChannelView that causes validator to fail - // PChannelView has 2 pchannels but current config's cluster has only 1 + // Return an assignment with an extra pchannel that causes validator to fail. + // The assignment has 2 pchannels but the current cluster configuration has only 1. b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - {Name: "other-1"}: channel.NewPChannelMeta("other-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1", "other-1"}, ReplicateConfiguration: &commonpb.ReplicateConfiguration{ Clusters: []*commonpb.MilvusCluster{ {ClusterId: "primary", Pchannels: []string{"primary-1"}, ConnectionParam: &commonpb.ConnectionParam{Uri: "http://primary:19530", Token: "primary"}}, @@ -1169,11 +1112,7 @@ func TestHandleForcePromoteNoCurrentConfig(t *testing.T) { b.EXPECT().Close().Return().Maybe() // Return assignment with nil ReplicateConfiguration b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, ReplicateConfiguration: nil, // No config exists }, nil).Maybe() balance.Register(b) @@ -1220,11 +1159,7 @@ func TestHandleForcePromoteClusterNotFound(t *testing.T) { b.EXPECT().Close().Return().Maybe() // Return config with different cluster ID (not "by-dev") b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, ReplicateConfiguration: &commonpb.ReplicateConfiguration{ Clusters: []*commonpb.MilvusCluster{ {ClusterId: "other-cluster", Pchannels: []string{"other-1"}, ConnectionParam: &commonpb.ConnectionParam{Uri: "http://other:19530", Token: "other"}}, @@ -1285,11 +1220,7 @@ func TestSecondValidateNonSameError(t *testing.T) { callCount++ if callCount <= 1 { return &balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1"}, }, nil } // Second call after lock: return error @@ -1389,12 +1320,7 @@ func TestForcePromoteMultiplePChannels(t *testing.T) { } // Multiple pchannels - by-dev-1 matches control channel, by-dev-2 does not b.EXPECT().GetLatestChannelAssignment().Return(&balancer.WatchChannelAssignmentsCallbackParam{ - PChannelView: &channel.PChannelView{ - Channels: map[channel.ChannelID]*channel.PChannelMeta{ - {Name: "by-dev-1"}: channel.NewPChannelMeta("by-dev-1", types.AccessModeRW), - {Name: "by-dev-2"}: channel.NewPChannelMeta("by-dev-2", types.AccessModeRW), - }, - }, + PChannels: []string{"by-dev-1", "by-dev-2"}, ReplicateConfiguration: currentReplicateConfig, }, nil).Maybe() balance.Register(b) diff --git a/internal/util/initcore/init_core.go b/internal/util/initcore/init_core.go index 350b4be5296..aaf35f3f13b 100644 --- a/internal/util/initcore/init_core.go +++ b/internal/util/initcore/init_core.go @@ -395,6 +395,7 @@ func InitTieredStorage(params *paramtable.ComponentParam) error { storageUsageTrackingEnabled := C.bool(params.QueryNodeCfg.StorageUsageTrackingEnabled.GetAsBool()) lazyManifestReaderEnabled := C.bool(params.QueryNodeCfg.TieredLazyManifestReaderEnabled.GetAsBool()) + lazyJsonStatsEnabled := C.bool(params.QueryNodeCfg.TieredLazyJsonStatsEnabled.GetAsBool()) evictionEnabled := C.bool(params.QueryNodeCfg.TieredEvictionEnabled.GetAsBool()) cacheTouchWindowMs := C.int64_t(params.QueryNodeCfg.TieredCacheTouchWindowMs.GetAsInt64()) backgroundEvictionEnabled := C.bool(params.QueryNodeCfg.TieredBackgroundEvictionEnabled.GetAsBool()) @@ -412,6 +413,7 @@ func InitTieredStorage(params *paramtable.ComponentParam) error { prefetchPoolThreads := C.uint32_t(hardware.GetCPUNum() * params.CommonCfg.LowPriorityThreadCoreCoefficient.GetAsInt()) C.SegcoreSetLazyManifestReaderEnabled(lazyManifestReaderEnabled) + C.SegcoreSetLazyJsonStatsEnabled(lazyJsonStatsEnabled) C.ConfigureTieredStorage(scalarFieldCacheWarmupPolicy, vectorFieldCacheWarmupPolicy, scalarIndexCacheWarmupPolicy, @@ -467,9 +469,11 @@ func UpdateTieredStorageConfig(params *paramtable.ComponentParam) error { warmupLoadingTimeoutMs := C.int64_t(params.QueryNodeCfg.TieredWarmupLoadingTimeoutMs.GetAsInt64()) storageUsageTrackingEnabled := C.bool(params.QueryNodeCfg.StorageUsageTrackingEnabled.GetAsBool()) lazyManifestReaderEnabled := C.bool(params.QueryNodeCfg.TieredLazyManifestReaderEnabled.GetAsBool()) + lazyJsonStatsEnabled := C.bool(params.QueryNodeCfg.TieredLazyJsonStatsEnabled.GetAsBool()) rejectRemoteVectorOutput := C.bool(params.QueryNodeCfg.TieredRejectRemoteVectorOutput.GetAsBool()) C.SegcoreSetLazyManifestReaderEnabled(lazyManifestReaderEnabled) + C.SegcoreSetLazyJsonStatsEnabled(lazyJsonStatsEnabled) C.UpdateTieredStorageConfig( loadingTimeoutMs, warmupLoadingTimeoutMs, @@ -795,6 +799,7 @@ func SetupCoreConfigChangelCallback() { paramtable.Get().QueryNodeCfg.TieredWarmupScalarIndex.RegisterCallback(updateTieredStorageConfigCallback) paramtable.Get().QueryNodeCfg.TieredWarmupVectorIndex.RegisterCallback(updateTieredStorageConfigCallback) paramtable.Get().QueryNodeCfg.TieredLazyManifestReaderEnabled.RegisterCallback(updateTieredStorageConfigCallback) + paramtable.Get().QueryNodeCfg.TieredLazyJsonStatsEnabled.RegisterCallback(updateTieredStorageConfigCallback) }) } diff --git a/pkg/util/paramtable/component_param.go b/pkg/util/paramtable/component_param.go index 903c65fc45a..ba53e398b46 100644 --- a/pkg/util/paramtable/component_param.go +++ b/pkg/util/paramtable/component_param.go @@ -3977,6 +3977,7 @@ type queryNodeConfig struct { StorageUsageTrackingEnabled ParamItem `refreshable:"true"` TieredRejectRemoteVectorOutput ParamItem `refreshable:"true"` TieredLazyManifestReaderEnabled ParamItem `refreshable:"true"` + TieredLazyJsonStatsEnabled ParamItem `refreshable:"true"` KnowhereScoreConsistency ParamItem `refreshable:"false"` @@ -4296,11 +4297,20 @@ Defaults to "sync".`, Key: "queryNode.segcore.tieredStorage.lazyManifestReaderEnabled", Version: "3.0.0", DefaultValue: "false", - Doc: "When enabled, Storage V3 manifest fields with warmup=disable and no load-time side effects defer projected ChunkReader, translator, and cache creation until first access.", + Doc: "When enabled, Storage V3 manifest fields with warmup=disable defer projected readers until first access.", Export: true, } p.TieredLazyManifestReaderEnabled.Init(base.mgr) + p.TieredLazyJsonStatsEnabled = ParamItem{ + Key: "queryNode.segcore.tieredStorage.lazyJsonStatsEnabled", + Version: "3.0.0", + DefaultValue: "true", + Doc: "When enabled, eligible Storage V3 JSON stats metadata initialization is deferred until first access.", + Export: true, + } + p.TieredLazyJsonStatsEnabled.Init(base.mgr) + p.TieredEvictionEnabled = ParamItem{ Key: "queryNode.segcore.tieredStorage.evictionEnabled", Version: "2.6.0", diff --git a/pkg/util/paramtable/component_param_test.go b/pkg/util/paramtable/component_param_test.go index aabbb64c47d..3e6d8ef737e 100644 --- a/pkg/util/paramtable/component_param_test.go +++ b/pkg/util/paramtable/component_param_test.go @@ -332,6 +332,23 @@ func TestComponentParam_LazyManifestReaderEnabled(t *testing.T) { assert.False(t, item.GetAsBool()) } +func TestComponentParam_LazyJsonStatsEnabled(t *testing.T) { + Init() + params := Get() + item := ¶ms.QueryNodeCfg.TieredLazyJsonStatsEnabled + t.Cleanup(func() { params.Reset(item.Key) }) + + assert.Equal(t, "queryNode.segcore.tieredStorage.lazyJsonStatsEnabled", item.Key) + assert.Equal(t, "true", item.DefaultValue) + assert.True(t, item.Export) + assert.True(t, item.GetAsBool()) + + assert.NoError(t, params.Save(item.Key, "false")) + assert.False(t, item.GetAsBool()) + assert.NoError(t, params.Save(item.Key, "true")) + assert.True(t, item.GetAsBool()) +} + func TestComponentParam_IDFLazyLoadSealedStats(t *testing.T) { Init() params := Get()