From 6bd2418df7687ff4bd0c40b4544171938e7286a9 Mon Sep 17 00:00:00 2001 From: sparknack Date: Wed, 15 Jul 2026 16:42:36 +0800 Subject: [PATCH] enhance: Load JsonStats shredding through manifest translator (#51116) issue: #51114 JsonStats shredding now loads through `ManifestGroupTranslator`/`ChunkReader`, but JsonStats files do not have the normal sealed-segment manifest metadata. The load path reconstructs the needed column-group view from the existing parquet footer/schema: column names, Arrow field ids, row counts, and file ordering. This keeps the current `meta.json` contract unchanged; parquet key-value metadata is only a fallback for old files that do not have the separate meta file. The main follow-up after switching to the manifest translator was projection behavior. Lazy loading now builds one projected reader per shredding column, so a predicate on one JSON path does not pull the whole group; when warmup is enabled it still uses the full column group so warmup can preload everything together. The added tests cover parquet files without packed field-list metadata, multi-file ordering, lazy single-column projection, and warmup full-group projection. --------- Signed-off-by: Shawn Wang --- .../expression/JsonContainsByStatsTest.cpp | 406 +++++++++++++++++- .../core/src/index/HybridScalarIndexTest.cpp | 40 +- .../src/index/json_stats/JsonKeyStats.cpp | 320 +++++++++++--- .../core/src/index/json_stats/JsonKeyStats.h | 2 + 4 files changed, 678 insertions(+), 90 deletions(-) diff --git a/internal/core/src/exec/expression/JsonContainsByStatsTest.cpp b/internal/core/src/exec/expression/JsonContainsByStatsTest.cpp index 4f2ae0e0fc..f09247b6f3 100644 --- a/internal/core/src/exec/expression/JsonContainsByStatsTest.cpp +++ b/internal/core/src/exec/expression/JsonContainsByStatsTest.cpp @@ -14,9 +14,11 @@ #include #include #include +#include #include #include #include +#include #include #include "NamedType/named_type_impl.hpp" @@ -37,9 +39,12 @@ #include "gtest/gtest.h" #include "index/IndexStats.h" #include "index/json_stats/JsonKeyStats.h" +#include "milvus-storage/common/constants.h" +#include "milvus-storage/common/metadata.h" #include "pb/common.pb.h" #include "pb/plan.pb.h" #include "pb/schema.pb.h" +#include "parquet/arrow/writer.h" #include "plan/PlanNode.h" #include "query/ExecPlanNodeVisitor.h" #include "segcore/ChunkedSegmentSealedImpl.h" @@ -60,6 +65,16 @@ using namespace milvus; using namespace milvus::index; +class JsonStatsProjectionTestAccessor { + public: + static bool + IsInMultiFieldColumnGroup(const JsonKeyStats& stats, + const std::string& field_name) { + return stats.shredding_columns_.at(field_name) + ->IsInMultiFieldColumnGroup(); + } +}; + namespace { bool @@ -87,17 +102,25 @@ MakeNullableJsonArray(const std::vector& json_strings, return std::static_pointer_cast(array); } -std::shared_ptr -BuildAndLoadJsonKeyStats(const std::vector& json_strings, - const milvus::FieldId json_fid, - const std::string& root_path, - int64_t collection_id, - int64_t partition_id, - int64_t segment_id, - int64_t field_id, - int64_t build_id, - int64_t version_id, - const std::vector* valid_data = nullptr) { +struct BuiltJsonStatsIndex { + storage::FileManagerContext ctx; + Config load_config; + std::vector index_files; + std::string stats_base_path; + milvus_storage::ArrowFileSystemPtr fs; +}; + +BuiltJsonStatsIndex +BuildJsonStatsIndex(const std::vector& json_strings, + const milvus::FieldId json_fid, + const std::string& root_path, + int64_t collection_id, + int64_t partition_id, + int64_t segment_id, + int64_t field_id, + int64_t build_id, + int64_t version_id, + const std::vector* valid_data = nullptr) { std::vector data; data.reserve(json_strings.size()); for (const auto& s : json_strings) { @@ -163,20 +186,254 @@ BuildAndLoadJsonKeyStats(const std::vector& json_strings, load_config["index_files"] = index_files; load_config[milvus::LOAD_PRIORITY] = milvus::proto::common::LoadPriority::HIGH; - load_config[STATS_BASE_PATH_KEY] = - storage::GenRemoteJsonStatsPathPrefix(chunk_manager, - build_id, - version_id, - collection_id, - partition_id, - segment_id, - field_id); + auto stats_base_path = storage::GenRemoteJsonStatsPathPrefix(chunk_manager, + build_id, + version_id, + collection_id, + partition_id, + segment_id, + field_id); + load_config[STATS_BASE_PATH_KEY] = stats_base_path; - auto reader = std::make_shared(ctx, true); - reader->Load(milvus::tracer::TraceContext{}, load_config); + return BuiltJsonStatsIndex{ctx, + std::move(load_config), + std::move(index_files), + std::move(stats_base_path), + std::move(fs)}; +} + +std::shared_ptr +LoadBuiltJsonStatsIndex(const BuiltJsonStatsIndex& built_index) { + auto reader = std::make_shared(built_index.ctx, true); + reader->Load(milvus::tracer::TraceContext{}, built_index.load_config); return reader; } +std::shared_ptr +BuildAndLoadJsonKeyStats(const std::vector& json_strings, + const milvus::FieldId json_fid, + const std::string& root_path, + int64_t collection_id, + int64_t partition_id, + int64_t segment_id, + int64_t field_id, + int64_t build_id, + int64_t version_id, + const std::vector* valid_data = nullptr) { + auto built_index = BuildJsonStatsIndex(json_strings, + json_fid, + root_path, + collection_id, + partition_id, + segment_id, + field_id, + build_id, + version_id, + valid_data); + return LoadBuiltJsonStatsIndex(built_index); +} + +std::string +FindFirstShreddingDataFile(const BuiltJsonStatsIndex& built_index) { + auto it = + std::find_if(built_index.index_files.begin(), + built_index.index_files.end(), + [](const std::string& file) { + return file.find(JSON_STATS_SHREDDING_DATA_PATH) != + std::string::npos; + }); + AssertInfo(it != built_index.index_files.end(), + "json stats index has no shredding data file"); + return built_index.stats_base_path + "/" + *it; +} + +std::string +MakeShreddingDataFile(const BuiltJsonStatsIndex& built_index, + int64_t column_group_id, + int64_t file_id) { + return fmt::format("{}/{}/{}/{}", + built_index.stats_base_path, + JSON_STATS_SHREDDING_DATA_PATH, + column_group_id, + file_id); +} + +std::string +MakeShreddingDataRelativeFile(int64_t column_group_id, int64_t file_id) { + return fmt::format( + "{}/{}/{}", JSON_STATS_SHREDDING_DATA_PATH, column_group_id, file_id); +} + +void +WriteShreddingParquetWithoutPackedFieldList( + const BuiltJsonStatsIndex& built_index, + int64_t column_group_id, + int64_t file_id, + const std::vector& values) { + auto path = MakeShreddingDataFile(built_index, column_group_id, file_id); + + arrow::Int64Builder value_builder; + AssertInfo(value_builder.AppendValues(values).ok(), + "failed to append json stats values"); + auto value_array = value_builder.Finish().ValueOrDie(); + + arrow::BinaryBuilder shared_builder; + for (size_t i = 0; i < values.size(); ++i) { + AssertInfo(shared_builder.AppendNull().ok(), + "failed to append shared json null"); + } + auto shared_array = shared_builder.Finish().ValueOrDie(); + + const auto value_field_id = START_JSON_STATS_FIELD_ID; + const auto shared_field_id = START_JSON_STATS_FIELD_ID + 1; + auto schema = arrow::schema({ + arrow::field( + JsonKey("/a", JSONType::INT64).ToColumnName(), + arrow::int64(), + true, + arrow::key_value_metadata({milvus_storage::ARROW_FIELD_ID_KEY}, + {std::to_string(value_field_id)})), + arrow::field( + JSON_KEY_STATS_SHARED_FIELD_NAME, + arrow::binary(), + true, + arrow::key_value_metadata({milvus_storage::ARROW_FIELD_ID_KEY}, + {std::to_string(shared_field_id)})), + }); + auto table = arrow::Table::Make(schema, {value_array, shared_array}); + + auto row_group_metadata = milvus_storage::RowGroupMetadataVector( + {milvus_storage::RowGroupMetadata(/*memory_size=*/128, + static_cast(values.size()), + /*row_offset=*/0)}); + auto file_metadata = + arrow::key_value_metadata({milvus_storage::ROW_GROUP_META_KEY, + milvus_storage::STORAGE_VERSION_KEY}, + {row_group_metadata.Serialize(), "2"}); + + auto output_result = built_index.fs->OpenOutputStream(path); + AssertInfo(output_result.ok(), + "failed to open parquet output {}: {}", + path, + output_result.status().ToString()); + auto output = output_result.ValueOrDie(); + + auto writer_result = parquet::arrow::FileWriter::Open( + *schema, arrow::default_memory_pool(), output); + AssertInfo(writer_result.ok(), + "failed to open parquet writer: {}", + writer_result.status().ToString()); + auto writer = std::move(writer_result).ValueOrDie(); + AssertInfo(writer->AddKeyValueMetadata(file_metadata).ok(), + "failed to add parquet metadata"); + AssertInfo(writer->WriteTable(*table, values.size()).ok(), + "failed to write parquet table"); + AssertInfo(writer->Close().ok(), "failed to close parquet writer"); + AssertInfo(output->Close().ok(), "failed to close parquet output"); +} + +void +OverwriteWithParquetMissingPackedFieldList( + const BuiltJsonStatsIndex& built_index, + const std::vector& values) { + (void)FindFirstShreddingDataFile(built_index); + WriteShreddingParquetWithoutPackedFieldList( + built_index, /*column_group_id=*/0, /*file_id=*/0, values); +} + +void +SetIndexFiles(BuiltJsonStatsIndex& built_index, + std::vector index_files) { + built_index.index_files = std::move(index_files); + built_index.load_config["index_files"] = built_index.index_files; +} + +TargetBitmap +ReadJsonStatsInt64Equal(JsonKeyStats& stats, + const std::string& field_name, + int64_t expected, + size_t size) { + TargetBitmap res(size); + TargetBitmap valid_res(size); + TargetBitmapView res_view(res); + TargetBitmapView valid_res_view(valid_res); + + auto func = [expected](const int64_t* data, + const bool* valid_data, + const int chunk_size, + TargetBitmapView res, + TargetBitmapView valid_res) { + for (int i = 0; i < chunk_size; ++i) { + valid_res[i] = valid_data[i]; + res[i] = valid_data[i] && data[i] == expected; + } + }; + + auto processed_size = stats.ExecutorForShreddingData( + nullptr, field_name, func, nullptr, res_view, valid_res_view); + AssertInfo(processed_size == size, + "processed json stats rows {} != {}", + processed_size, + size); + return res; +} + +void +AssertJsonStatsProjectionMode(const std::string& warmup_policy, + int64_t id_offset, + bool expect_multi_field_group) { + auto schema = std::make_shared(); + auto json_fid = schema->AddDebugField("json", DataType::JSON); + + std::vector json_raw_data = { + R"({"a": 1, "b": 10})", + R"({"a": 2, "b": 20})", + R"({"a": 1, "b": 30})", + }; + + const int64_t collection_id = 1200 + id_offset; + const int64_t partition_id = 2200 + id_offset; + const int64_t segment_id = 3200 + id_offset; + const int64_t field_id = json_fid.get(); + const int64_t build_id = 5200 + id_offset; + const int64_t version_id = 1; + const std::string root_path = TestLocalPath; + + auto built_index = BuildJsonStatsIndex(json_raw_data, + json_fid, + root_path, + collection_id, + partition_id, + segment_id, + field_id, + build_id, + version_id); + built_index.load_config[milvus::index::WARMUP] = warmup_policy; + ASSERT_TRUE(built_index.load_config.contains(milvus::index::WARMUP)); + ASSERT_EQ( + built_index.load_config.at(milvus::index::WARMUP).get(), + warmup_policy); + auto stats = LoadBuiltJsonStatsIndex(built_index); + + auto a_field = stats->GetShreddingField("/a", JSONType::INT64); + auto b_field = stats->GetShreddingField("/b", JSONType::INT64); + ASSERT_FALSE(a_field.empty()); + ASSERT_FALSE(b_field.empty()); + + auto a_result = + ReadJsonStatsInt64Equal(*stats, a_field, /*expected=*/1, /*size=*/3); + EXPECT_TRUE(a_result[0]); + EXPECT_FALSE(a_result[1]); + EXPECT_TRUE(a_result[2]); + + EXPECT_EQ(JsonStatsProjectionTestAccessor::IsInMultiFieldColumnGroup( + *stats, a_field), + expect_multi_field_group); + EXPECT_EQ(JsonStatsProjectionTestAccessor::IsInMultiFieldColumnGroup( + *stats, b_field), + expect_multi_field_group); +} + } // namespace TEST(JsonContainsByStatsTest, BasicContainsAnyOnArray) { @@ -282,6 +539,113 @@ TEST(JsonContainsByStatsTest, BasicContainsAnyOnArray) { } } +TEST(JsonStatsAsyncLoadTest, LoadsShreddingParquetWithoutPackedFieldList) { + auto schema = std::make_shared(); + auto json_fid = schema->AddDebugField("json", DataType::JSON); + + std::vector json_raw_data = { + R"({"a": 1})", + R"({"a": 2})", + R"({"a": 1})", + R"({"a": 3})", + }; + + const int64_t collection_id = 1201; + const int64_t partition_id = 2201; + const int64_t segment_id = 3201; + const int64_t field_id = json_fid.get(); + const int64_t build_id = 5201; + const int64_t version_id = 1; + const std::string root_path = TestLocalPath; + + auto built_index = BuildJsonStatsIndex(json_raw_data, + json_fid, + root_path, + collection_id, + partition_id, + segment_id, + field_id, + build_id, + version_id); + OverwriteWithParquetMissingPackedFieldList( + built_index, std::vector{1, 2, 1, 3}); + + auto stats = LoadBuiltJsonStatsIndex(built_index); + auto result = ReadJsonStatsInt64Equal( + *stats, JsonKey("/a", JSONType::INT64).ToColumnName(), 1, 4); + + EXPECT_TRUE(result[0]); + EXPECT_FALSE(result[1]); + EXPECT_TRUE(result[2]); + EXPECT_FALSE(result[3]); + EXPECT_EQ(result.count(), 2); +} + +TEST(JsonStatsAsyncLoadTest, LoadsMultipleShreddingParquetFilesInFileIdOrder) { + auto schema = std::make_shared(); + auto json_fid = schema->AddDebugField("json", DataType::JSON); + + std::vector json_raw_data = { + R"({"a": 1})", + R"({"a": 2})", + R"({"a": 3})", + R"({"a": 1})", + R"({"a": 4})", + }; + + const int64_t collection_id = 1202; + const int64_t partition_id = 2202; + const int64_t segment_id = 3202; + const int64_t field_id = json_fid.get(); + const int64_t build_id = 5202; + const int64_t version_id = 1; + const std::string root_path = TestLocalPath; + + auto built_index = BuildJsonStatsIndex(json_raw_data, + json_fid, + root_path, + collection_id, + partition_id, + segment_id, + field_id, + build_id, + version_id); + WriteShreddingParquetWithoutPackedFieldList( + built_index, /*column_group_id=*/0, /*file_id=*/0, {1, 2}); + WriteShreddingParquetWithoutPackedFieldList( + built_index, /*column_group_id=*/0, /*file_id=*/1, {3, 1, 4}); + + std::vector shuffled_index_files{ + MakeShreddingDataRelativeFile(/*column_group_id=*/0, /*file_id=*/1)}; + shuffled_index_files.insert(shuffled_index_files.end(), + built_index.index_files.begin(), + built_index.index_files.end()); + SetIndexFiles(built_index, std::move(shuffled_index_files)); + + auto stats = LoadBuiltJsonStatsIndex(built_index); + auto result = ReadJsonStatsInt64Equal( + *stats, JsonKey("/a", JSONType::INT64).ToColumnName(), 1, 5); + + EXPECT_TRUE(result[0]); + EXPECT_FALSE(result[1]); + EXPECT_FALSE(result[2]); + EXPECT_TRUE(result[3]); + EXPECT_FALSE(result[4]); + EXPECT_EQ(result.count(), 2); +} + +TEST(JsonStatsAsyncLoadTest, UsesSingleColumnProjectionWithoutWarmup) { + AssertJsonStatsProjectionMode("disable", + /*id_offset=*/3, + /*expect_multi_field_group=*/false); +} + +TEST(JsonStatsAsyncLoadTest, UsesFullColumnGroupProjectionWithWarmup) { + AssertJsonStatsProjectionMode("sync", + /*id_offset=*/4, + /*expect_multi_field_group=*/true); +} + TEST(JsonStatsUnaryRangeTest, NotEqualKeepsJsonPathUnknownsAndMasksFieldNull) { auto schema = std::make_shared(); auto json_fid = schema->AddDebugField("json", DataType::JSON, true); diff --git a/internal/core/src/index/HybridScalarIndexTest.cpp b/internal/core/src/index/HybridScalarIndexTest.cpp index 62b9b83b98..f138952b1e 100644 --- a/internal/core/src/index/HybridScalarIndexTest.cpp +++ b/internal/core/src/index/HybridScalarIndexTest.cpp @@ -19,12 +19,14 @@ #include #include #include +#include #include #include #include #include #include #include +#include #include #include #include @@ -284,8 +286,13 @@ class HybridIndexTestV1 : public testing::Test { index_version_); } + void + TearDown() override { + CleanupLocalRoot(/*report_error=*/true); + } + virtual ~HybridIndexTestV1() override { - boost::filesystem::remove_all(chunk_manager_->GetRootPath()); + CleanupLocalRoot(/*report_error=*/false); } public: @@ -568,6 +575,37 @@ class HybridIndexTestV1 : public testing::Test { bool has_lack_binlog_row_{false}; size_t lack_binlog_row_{100}; std::string hybrid_high_cardinality_index_type_; + + private: + void + CleanupLocalRoot(bool report_error) { + index_.reset(); + if (chunk_manager_ == nullptr) { + return; + } + + const auto root_path = chunk_manager_->GetRootPath(); + boost::system::error_code ec; + for (int attempt = 0; attempt < 3; ++attempt) { + ec.clear(); + boost::filesystem::remove_all(root_path, ec); + if (!ec) { + return; + } + + boost::system::error_code exists_ec; + if (!boost::filesystem::exists(root_path, exists_ec) && + !exists_ec) { + return; + } + std::this_thread::sleep_for(std::chrono::milliseconds(50)); + } + + if (report_error) { + ADD_FAILURE() << "failed to remove hybrid index test root " + << root_path << ": " << ec.message(); + } + } }; TYPED_TEST_SUITE_P(HybridIndexTestV1); diff --git a/internal/core/src/index/json_stats/JsonKeyStats.cpp b/internal/core/src/index/json_stats/JsonKeyStats.cpp index a8f0ccd370..147ed75fa9 100644 --- a/internal/core/src/index/json_stats/JsonKeyStats.cpp +++ b/internal/core/src/index/json_stats/JsonKeyStats.cpp @@ -47,8 +47,11 @@ #include "milvus-storage/common/config.h" #include "milvus-storage/common/constants.h" #include "milvus-storage/common/metadata.h" +#include "milvus-storage/column_groups.h" #include "milvus-storage/filesystem/fs.h" #include "milvus-storage/format/parquet/file_reader.h" +#include "milvus-storage/format/parquet/parquet_format_reader.h" +#include "milvus-storage/reader.h" #include "mmap/ChunkedColumnGroup.h" #include "mmap/Types.h" #include "nlohmann/detail/iterators/iteration_proxy.hpp" @@ -56,7 +59,7 @@ #include "parquet/metadata.h" #include "segcore/storagev1translator/BsonInvertedIndexTranslator.h" #include "segcore/ChunkedSegmentSealedImpl.h" -#include "segcore/storagev2translator/GroupChunkTranslator.h" +#include "segcore/storagev2translator/ManifestGroupTranslator.h" #include "segcore/Utils.h" #include "storage/DiskFileManagerImpl.h" #include "storage/FileManager.h" @@ -69,9 +72,97 @@ #include "storage/ThreadPools.h" #include "storage/Types.h" #include "storage/Util.h" +#include "storage/loon_ffi/property_singleton.h" namespace milvus::index { +namespace { + +struct JsonStatsParquetMetadata { + std::shared_ptr schema; + int64_t num_rows; +}; + +milvus_storage::api::Properties +GetJsonStatsReadProperties() { + auto properties = + storage::LoonFFIPropertiesSingleton::GetInstance().GetProperties(); + if (properties == nullptr) { + return {}; + } + return *properties; +} + +std::string +NoopParquetKeyRetriever(const std::string&) { + return {}; +} + +JsonStatsParquetMetadata +ReadJsonStatsParquetMetadata(const std::string& file) { + auto fs = milvus::segcore::GetDefaultArrowFileSystem(); + auto properties = GetJsonStatsReadProperties(); + milvus_storage::parquet::ParquetFormatReader reader( + fs, file, properties, {}, NoopParquetKeyRetriever); + + auto open_status = reader.open(); + AssertInfo(open_status.ok(), + "[JsonStats] failed to open parquet metadata reader for {}: {}", + file, + open_status.ToString()); + + auto row_group_result = reader.get_row_group_infos(); + AssertInfo(row_group_result.ok(), + "[JsonStats] failed to read parquet row groups for {}: {}", + file, + row_group_result.status().ToString()); + auto row_groups = row_group_result.ValueOrDie(); + + int64_t num_rows = 0; + for (const auto& row_group : row_groups) { + num_rows += + static_cast(row_group.end_offset - row_group.start_offset); + } + + auto schema = reader.get_schema(); + AssertInfo(schema != nullptr, + "[JsonStats] failed to read parquet schema for {}", + file); + return JsonStatsParquetMetadata{std::move(schema), num_rows}; +} + +FieldId +GetJsonStatsFieldIdFromArrowField(const std::shared_ptr& field) { + const auto& metadata = field->metadata(); + AssertInfo(metadata != nullptr && + metadata->Contains(milvus_storage::ARROW_FIELD_ID_KEY), + "json stats field id not found in metadata for field {}", + field->name()); + auto result = metadata->Get(milvus_storage::ARROW_FIELD_ID_KEY); + AssertInfo(result.ok(), + "failed to get json stats field id from metadata for field {}: " + "{}", + field->name(), + result.status().ToString()); + return FieldId(std::stoll(result.ValueOrDie())); +} + +std::pair, std::vector> +GetJsonStatsFieldsFromSchema(const std::shared_ptr& schema) { + std::vector field_ids; + std::vector field_names; + field_ids.reserve(schema->num_fields()); + field_names.reserve(schema->num_fields()); + + for (const auto& field : schema->fields()) { + field_ids.push_back(GetJsonStatsFieldIdFromArrowField(field)); + field_names.push_back(field->name()); + } + return {std::move(field_ids), std::move(field_names)}; +} + +} // namespace + JsonKeyStats::JsonKeyStats(const storage::FileManagerContext& ctx, bool is_load, int64_t json_stats_max_shredding_columns, @@ -793,13 +884,8 @@ JsonKeyStats::BuildWithFieldData(const std::vector& field_datas, void JsonKeyStats::GetColumnSchemaFromParquet(int64_t column_group_id, const std::string& file) { - auto fs = milvus::segcore::GetDefaultArrowFileSystem(); - auto result = milvus_storage::FileRowGroupReader::Make(fs, file); - AssertInfo(result.ok(), - "[StorageV2] Failed to create file row group reader: {}", - result.status().ToString()); - auto file_reader = result.ValueOrDie(); - std::shared_ptr file_schema = file_reader->schema(); + auto parquet_metadata = ReadJsonStatsParquetMetadata(file); + std::shared_ptr file_schema = parquet_metadata.schema; LOG_DEBUG("get column schema: [{}] for segment {}", file_schema->ToString(true), segment_id_); @@ -963,40 +1049,27 @@ JsonKeyStats::LoadColumnGroup(int64_t column_group_id, remote_prefix, column_group_id, file_id)); } - auto fs = milvus::segcore::GetDefaultArrowFileSystem(); - auto result = milvus_storage::FileRowGroupReader::Make(fs, files[0]); - AssertInfo(result.ok(), - "[StorageV2] Failed to create file row group reader: {}", - result.status().ToString()); - auto file_reader = result.ValueOrDie(); - std::shared_ptr metadata = - file_reader->file_metadata(); - milvus_storage::FieldIDList field_id_list = - metadata->GetGroupFieldIDList().GetFieldIDList(column_group_id); - std::vector milvus_field_ids; - for (int i = 0; i < field_id_list.size(); ++i) { - milvus_field_ids.push_back(FieldId(field_id_list.Get(i))); - } + auto first_file_metadata = ReadJsonStatsParquetMetadata(files[0]); + auto [milvus_field_ids, column_names] = + GetJsonStatsFieldsFromSchema(first_file_metadata.schema); - // Fetch row group metadata from all files in parallel using HIGH POOL - // to avoid blocking the caller thread with serial S3 I/O + std::vector file_num_rows; + file_num_rows.reserve(files.size()); + file_num_rows.push_back(first_file_metadata.num_rows); + num_rows += first_file_metadata.num_rows; + + // Fetch row group metadata from remaining files in parallel using HIGH POOL + // to avoid blocking the caller thread with serial S3 I/O. auto& pool = ThreadPools::GetThreadPool(milvus::ThreadPoolPriority::HIGH); std::vector> futures; - futures.reserve(files.size()); - for (const auto& file : files) { - futures.push_back(pool.Submit([&fs, file]() { - auto result = milvus_storage::FileRowGroupReader::Make(fs, file); - AssertInfo(result.ok(), - "[StorageV2] Failed to create file row group reader: " + - result.status().ToString()); - auto reader = result.ValueOrDie(); - auto row_group_meta_vector = - reader->file_metadata()->GetRowGroupMetadataVector(); - return static_cast(row_group_meta_vector.row_num()); - })); + futures.reserve(files.size() - 1); + for (size_t i = 1; i < files.size(); ++i) { + const auto& file = files[i]; + futures.push_back(pool.Submit( + [file]() { return ReadJsonStatsParquetMetadata(file).num_rows; })); } // Ensure all futures are awaited even if one throws, to prevent - // use-after-free on captured references (&fs) in background tasks. + // use-after-free on captured references in background tasks. auto futures_guard = folly::makeGuard([&futures]() { for (auto& f : futures) { if (f.valid()) { @@ -1008,7 +1081,9 @@ JsonKeyStats::LoadColumnGroup(int64_t column_group_id, } }); for (auto& f : futures) { - num_rows += f.get(); + auto file_rows = f.get(); + file_num_rows.push_back(file_rows); + num_rows += file_rows; } if (num_rows_ == 0) { @@ -1019,8 +1094,6 @@ JsonKeyStats::LoadColumnGroup(int64_t column_group_id, segment_id_); auto enable_mmap = !mmap_filepath_.empty(); - auto column_group_info = FieldDataInfo( - column_group_id, field_id_, num_rows, mmap_filepath_, false, shard_); LOG_INFO( "loads column group {} with num_rows {} for segment " "{}", @@ -1029,44 +1102,156 @@ JsonKeyStats::LoadColumnGroup(int64_t column_group_id, segment_id_); std::unordered_map field_meta_map; - for (const auto& inner_field_id : milvus_field_ids) { - auto field_name = field_id_to_name_map_[inner_field_id.get()]; + for (size_t i = 0; i < milvus_field_ids.size(); ++i) { + const auto& inner_field_id = milvus_field_ids[i]; + auto field_name_it = field_id_to_name_map_.find(inner_field_id.get()); + AssertInfo(field_name_it != field_id_to_name_map_.end(), + "field id {} not found in json stats field map for " + "segment {}", + inner_field_id.get(), + segment_id_); + auto field_name = field_name_it->second; FieldMeta field_meta( FieldName(field_name), inner_field_id, field_id_, GetPrimitiveDataType(shred_field_data_type_map_[field_name]), true, - std::nullopt); + std::nullopt, + column_names[i]); field_meta_map.insert(std::make_pair(FieldId(inner_field_id.get()), std::move(field_meta))); } auto& mmap_config = storage::MmapManager::GetInstance().GetMmapConfig(); - auto group_chunk_metadata = milvus::segcore::LoadGroupChunkMetadata( - files, {}, fmt::format("seg_{}_jks_{}", segment_id_, field_id_)); + auto column_group = std::make_shared(); + column_group->columns = column_names; + column_group->format = LOON_FORMAT_PARQUET; + column_group->files.reserve(files.size()); + for (size_t i = 0; i < files.size(); ++i) { + column_group->files.push_back(milvus_storage::api::ColumnGroupFile{ + .path = files[i], + .start_index = 0, + .end_index = file_num_rows[i], + .properties = {}, + }); + } + auto column_groups = std::make_shared(); + column_groups->push_back(std::move(column_group)); + auto properties = GetJsonStatsReadProperties(); + auto resolved_warmup_policy = + milvus::segcore::getCacheWarmupPolicy(warmup_policy, + /*is_vector=*/false, + /*is_index=*/false, + /*in_load_list=*/true); + auto eager_load = + resolved_warmup_policy != CacheWarmupPolicy::CacheWarmupPolicy_Disable; - auto translator = std::make_unique< - milvus::segcore::storagev2translator::GroupChunkTranslator>( - segment_id_, - GroupChunkType::JSON_KEY_STATS, - field_meta_map, - column_group_info, - std::move(files), - std::move(group_chunk_metadata.row_group_meta_list), - enable_mmap, - mmap_config.GetMmapPopulate(), - milvus_field_ids.size(), - load_priority_, - warmup_policy); + if (eager_load) { + auto needed_columns = + std::make_shared>(column_names); + auto reader = milvus_storage::api::Reader::create( + column_groups, nullptr, needed_columns, properties); + auto chunk_reader_result = reader->get_chunk_reader(0, needed_columns); + AssertInfo(chunk_reader_result.ok(), + "[JsonStats] failed to create chunk reader for column group " + "{} segment {}: {}", + column_group_id, + segment_id_, + chunk_reader_result.status().ToString()); + auto chunk_reader_unique = std::move(chunk_reader_result).ValueOrDie(); + std::shared_ptr chunk_reader( + std::move(chunk_reader_unique)); - auto chunked_column_group = - std::make_shared(std::move(translator)); + auto translator = std::make_unique< + milvus::segcore::storagev2translator::ManifestGroupTranslator>( + segment_id_, + GroupChunkType::JSON_KEY_STATS, + column_group_id, + std::move(chunk_reader), + field_meta_map, + enable_mmap, + mmap_config.GetMmapPopulate(), + mmap_filepath_, + milvus_field_ids.size(), + load_priority_, + /*eager_load=*/true, + warmup_policy, + fmt::format("jks_{}", field_id_), + /*fallback_bytes_per_row=*/0, + shard_); + + auto chunked_column_group = + std::make_shared(std::move(translator)); + + for (const auto& inner_field_id : milvus_field_ids) { + auto field_meta = field_meta_map.at(inner_field_id); + auto column = std::make_shared( + chunked_column_group, inner_field_id, field_meta); + + LOG_DEBUG( + "add shredding column: {}, inner_field_id:{}, for json field " + "{} segment " + "{}", + field_meta.get_name().get(), + inner_field_id.get(), + field_id_, + segment_id_); + shredding_columns_[field_meta.get_name().get()] = column; + } + shared_column_ = shredding_columns_.at(shared_column_field_name_); + return; + } + + // Lazy JSON stats columns are loaded through per-column projected readers, + // same as lazy storage-v2 column-group entries. This avoids co-loading + // sibling JSON paths when a query touches only one shredding column. + auto all_columns = std::make_shared>(column_names); + auto reader = milvus_storage::api::Reader::create( + column_groups, nullptr, all_columns, properties); + for (size_t i = 0; i < milvus_field_ids.size(); ++i) { + const auto& inner_field_id = milvus_field_ids[i]; + const auto& column_name = column_names[i]; + auto needed_columns = + std::make_shared>(std::vector{ + column_name, + }); + auto chunk_reader_result = reader->get_chunk_reader(0, needed_columns); + AssertInfo(chunk_reader_result.ok(), + "[JsonStats] failed to create projected chunk reader for " + "column group {} column {} segment {}: {}", + column_group_id, + column_name, + segment_id_, + chunk_reader_result.status().ToString()); + auto chunk_reader_unique = std::move(chunk_reader_result).ValueOrDie(); + std::shared_ptr chunk_reader( + std::move(chunk_reader_unique)); - // Create ProxyChunkColumn for each field in this column group - for (const auto& inner_field_id : milvus_field_ids) { auto field_meta = field_meta_map.at(inner_field_id); + std::unordered_map projected_field_meta_map; + projected_field_meta_map.emplace(inner_field_id, field_meta); + auto translator = std::make_unique< + milvus::segcore::storagev2translator::ManifestGroupTranslator>( + segment_id_, + GroupChunkType::JSON_KEY_STATS, + column_group_id, + std::move(chunk_reader), + projected_field_meta_map, + enable_mmap, + mmap_config.GetMmapPopulate(), + mmap_filepath_, + /*num_fields=*/1, + load_priority_, + eager_load, + warmup_policy, + fmt::format("jks_{}_{}", field_id_, inner_field_id.get()), + /*fallback_bytes_per_row=*/0, + shard_); + + auto chunked_column_group = + std::make_shared(std::move(translator)); auto column = std::make_shared( chunked_column_group, inner_field_id, field_meta); @@ -1163,6 +1348,8 @@ JsonKeyStats::Load(milvus::tracer::TraceContext ctx, const Config& config) { static_cast(load_priority_)); shard_ = GetValueFromConfig(config, JSON_STATS_CACHE_SHARD_KEY) .value_or(""); + auto warmup_policy = + GetValueFromConfig(config, WARMUP).value_or(""); auto index_files = GetValueFromConfig>(config, "index_files"); @@ -1215,17 +1402,14 @@ JsonKeyStats::Load(milvus::tracer::TraceContext ctx, const Config& config) { } // load shredding data (files are already absolute paths) - LoadShreddingData(shredding_data_files, - config.contains(WARMUP) ? config.at(WARMUP) : ""); + LoadShreddingData(shredding_data_files, warmup_policy); auto index_size = GetValueFromConfig(config, milvus::index::INDEX_SIZE) .value_or(0); // load shared key index (files are already absolute paths) - LoadSharedKeyIndex(shared_key_index_files, - enable_mmap, - index_size, - config.contains(WARMUP) ? config.at(WARMUP) : ""); + LoadSharedKeyIndex( + shared_key_index_files, enable_mmap, index_size, warmup_policy); } IndexStatsPtr diff --git a/internal/core/src/index/json_stats/JsonKeyStats.h b/internal/core/src/index/json_stats/JsonKeyStats.h index ff35f98ea9..ccca5dae57 100644 --- a/internal/core/src/index/json_stats/JsonKeyStats.h +++ b/internal/core/src/index/json_stats/JsonKeyStats.h @@ -70,6 +70,7 @@ class CollectSingleJsonStatsInfoAccessor; // Forward declaration of test accessor in global namespace for friend declaration class TraverseJsonForBuildStatsAccessor; +class JsonStatsProjectionTestAccessor; namespace milvus::index { class JsonKeyStats : public ScalarIndex { @@ -720,6 +721,7 @@ class JsonKeyStats : public ScalarIndex { // Friend accessor for unit tests to call private methods safely. friend class ::TraverseJsonForBuildStatsAccessor; friend class ::CollectSingleJsonStatsInfoAccessor; + friend class ::JsonStatsProjectionTestAccessor; }; } // namespace milvus::index