airborne12 commented on code in PR #67977: URL: https://github.com/apache/doris/pull/67977#discussion_r4214413042
########## be/src/storage/index/index_disk_usage.cpp: ########## @@ -0,0 +1,413 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include "storage/index/index_disk_usage.h" + +#include <algorithm> +#include <map> +#include <memory> +#include <tuple> +#include <utility> + +#include "common/cast_set.h" +#include "common/check.h" +#include "io/io_common.h" +#include "storage/index/ann/ann_index_files.h" +#include "storage/index/index_file_reader.h" +#include "storage/index/inverted/inverted_index_desc.h" +#include "storage/index/snii/format/dict_entry.h" +#include "storage/index/snii/format/format_constants.h" +#include "storage/index/snii/reader/logical_index_reader.h" +#include "storage/olap_common.h" +#include "storage/rowset/beta_rowset.h" +#include "storage/rowset/rowset.h" +#include "storage/rowset/rowset_meta.h" + +namespace doris::segment_v2 { + +namespace { + +bool is_bkd_file(std::string_view name) { + return name == InvertedIndexDescriptor::get_temporary_bkd_index_data_file_name() || + name == InvertedIndexDescriptor::get_temporary_bkd_index_meta_file_name() || + name == InvertedIndexDescriptor::get_temporary_bkd_index_file_name(); +} + +bool is_ann_file(std::string_view name) { + return name == faiss_index_fila_name || name == faiss_ivfdata_file_name; +} + +bool is_wanted(const IndexDiskUsageOptions& options, int64_t index_id) { + return options.index_ids.empty() || options.index_ids.contains(index_id); +} + +Status check_cancelled(const IndexDiskUsageOptions& options) { + return options.check_cancelled ? options.check_cancelled() : Status::OK(); +} + +// A segment may have no index file on purpose, for example when every ANN index skipped a segment +// too small to train, or when a legacy table skipped writing indexes on load. Such a file holds +// no index bytes, while any other error is still reported. +bool is_absent_index_file(const Status& status) { + return status.is<ErrorCode::INVERTED_INDEX_FILE_NOT_FOUND>() || + status.is<ErrorCode::INVERTED_INDEX_BYPASS>() || status.is<ErrorCode::NOT_FOUND>(); +} + +// Returns the V1 index file size persisted in the rowset meta, or -1 when it is not recorded. +int64_t persisted_v1_file_size(const InvertedIndexFileInfo& file_info, const TabletIndex& index) { + for (const auto& index_info : file_info.index_info()) { + if (index_info.index_id() == index.index_id() && + index_info.index_suffix() == index.get_index_suffix()) { + return index_info.index_file_size() > 0 ? index_info.index_file_size() : -1; + } + } + return -1; +} + +// Classifies every sub-file of a CLucene directory into `record` and adds their total length to +// `files_bytes`. +Status add_directory_files(const lucene::store::Directory& dir, IndexDiskUsageRecord* record, + int64_t* files_bytes) { + try { + std::vector<std::string> names; + if (!dir.list(&names)) { + return Status::Error<ErrorCode::INVERTED_INDEX_CLUCENE_ERROR>( + "failed to list inverted index sub-files"); + } + for (const auto& name : names) { + const int64_t length = dir.fileLength(name.c_str()); + classify_clucene_file(name, length, record); + *files_bytes += length; + } + } catch (CLuceneError& e) { + return Status::Error<ErrorCode::INVERTED_INDEX_CLUCENE_ERROR>( + "failed to read inverted index sub-files: {}", e.what()); + } + return Status::OK(); +} + +// Sums the position bytes of the dictionary entries whose postings live in the posting region. +// Inline postings stay in the dictionary region and are not counted. +Status sum_snii_position_bytes(const IndexFileReader& reader, uint64_t index_id, + std::string_view suffix, const IndexDiskUsageOptions& options, + int64_t* position_bytes) { + // A full dictionary scan should not evict blocks that queries keep in the file cache. + io::IOContext io_ctx; + io_ctx.is_disposable = true; + io_ctx.is_inverted_index = true; + auto logical = DORIS_TRY(reader.open_snii_logical_index( + index_id, suffix, &io_ctx, snii::reader::LogicalIndexOpenMode::kCompaction)); + snii_doris::DorisSniiFileReader::ScopedIOContext io_context_scope(&io_ctx); + std::vector<snii::format::DictEntry> entries; + for (uint32_t block = 0; block < logical->n_dict_blocks(); ++block) { + RETURN_IF_ERROR(check_cancelled(options)); + uint64_t frq_base = 0; + uint64_t prx_base = 0; + RETURN_IF_ERROR(logical->decode_dict_block(block, &entries, &frq_base, &prx_base)); + for (const auto& entry : entries) { + if (entry.kind == snii::format::DictEntryKind::kPodRef) { + *position_bytes += cast_set<int64_t>(entry.prx_len); + } + } + } + return Status::OK(); +} + +int64_t merge_component(int64_t lhs, int64_t rhs) { + return lhs < 0 || rhs < 0 ? -1 : lhs + rhs; +} + +} // namespace + +std::vector<IndexDiskUsageRow> aggregate_index_disk_usage(std::vector<IndexDiskUsageRow> rows, + IndexDiskUsageLevel level) { + if (level == IndexDiskUsageLevel::kSegment) { + return rows; + } + using Key = std::tuple<std::string, int64_t, std::string, int, int>; + std::map<Key, size_t> positions; + std::vector<IndexDiskUsageRow> merged; + for (auto& row : rows) { + row.segment_id = -1; + if (level == IndexDiskUsageLevel::kTablet) { + row.rowset_id.clear(); + } + Key key {row.rowset_id, row.record.index_id, row.record.index_suffix, + static_cast<int>(row.format), static_cast<int>(row.record.structure)}; + auto [it, inserted] = positions.emplace(std::move(key), merged.size()); + if (inserted) { + merged.push_back(std::move(row)); + continue; + } + IndexDiskUsageRow& target = merged[it->second]; + target.segment_count += row.segment_count; + target.row_count += row.row_count; + IndexDiskUsageRecord& dst = target.record; + const IndexDiskUsageRecord& src = row.record; + dst.total_bytes += src.total_bytes; + dst.dict_bytes = merge_component(dst.dict_bytes, src.dict_bytes); + dst.posting_bytes = merge_component(dst.posting_bytes, src.posting_bytes); + dst.position_bytes = merge_component(dst.position_bytes, src.position_bytes); + dst.stats_bytes = merge_component(dst.stats_bytes, src.stats_bytes); + dst.other_bytes = merge_component(dst.other_bytes, src.other_bytes); + } + return merged; +} + +void classify_clucene_file(std::string_view name, int64_t length, IndexDiskUsageRecord* record) { + record->total_bytes += length; + if (is_bkd_file(name)) { + record->structure = IndexDiskUsageStructure::kBkd; + return; + } + if (is_ann_file(name)) { + record->structure = IndexDiskUsageStructure::kAnn; + return; + } + const size_t dot = name.rfind('.'); + const std::string_view extension = + dot == std::string_view::npos ? std::string_view() : name.substr(dot + 1); + if (extension == "tis" || extension == "tii") { + record->dict_bytes += length; + } else if (extension == "frq") { + record->posting_bytes += length; + } else if (extension == "prx") { + record->position_bytes += length; + } else if (extension == "nrm") { + record->stats_bytes += length; + } else { + record->other_bytes += length; + } +} + +IndexDiskUsageCollector::IndexDiskUsageCollector(io::FileSystemSPtr fs, + std::string index_path_prefix, + TabletSchemaSPtr schema, + InvertedIndexStorageFormatPB format, + int64_t tablet_id, + InvertedIndexFileInfo index_file_info) + : _fs(std::move(fs)), + _index_path_prefix(std::move(index_path_prefix)), + _schema(std::move(schema)), + _format(format), + _tablet_id(tablet_id), + _index_file_info(std::move(index_file_info)) {} + +Status IndexDiskUsageCollector::collect(const IndexDiskUsageOptions& options, + std::vector<IndexDiskUsageRecord>* out) { + DORIS_CHECK(out != nullptr); + switch (_format) { + case InvertedIndexStorageFormatPB::V1: + return _collect_v1(options, out); + case InvertedIndexStorageFormatPB::V2: + case InvertedIndexStorageFormatPB::V3: + return _collect_compound(options, out); + case InvertedIndexStorageFormatPB::SNII: + return _collect_snii(options, out); + default: + return Status::NotSupported("index disk usage does not support inverted index format {}", + InvertedIndexStorageFormatPB_Name(_format)); + } +} + +Status IndexDiskUsageCollector::_collect_v1(const IndexDiskUsageOptions& options, + std::vector<IndexDiskUsageRecord>* out) { + IndexFileReader reader(_fs, _index_path_prefix, _format, _index_file_info, _tablet_id); + RETURN_IF_ERROR(reader.init()); + for (const TabletIndex* index : _schema->inverted_indexes()) { Review Comment: The ANN half of this finding is fixed in 0b579993e1a. For VARIANT paths I did not change the enumeration, and added a regression test instead, for these reasons. Rowsets written since `index_info` exists are covered by it, because it records every `(index_id, suffix)` file that was written (`CollectV1VariantPathFilesFromFileInfo`). Rowsets that predate `index_info` use the fallback over their own schema, and the writers of that era persisted one schema index per extracted path: the rowset schema is merged with `get_least_common_schema`, whose `inherit_column_attributes(schema)` appends the parent index with the path as `index_suffix` (branch-2.1 and branch-3.0 `schema_util.cpp`, and `variant_util.cpp` on master). The rest of the BE locates V1 files through the schema alone: `BetaRowset::remove()` and `copy_files_to()` walk `_schema->inverted_indexs(*column)` for every column, extracted ones included. The other writers that leave a rowset without `index_info` do not add extra V1 files for VARIANT either: segcompaction is disabled for schemas with VARIANT columns (`BetaRowsetWriter::_segcompaction_if_necessary`), and `IndexBuilder` skips VAR IANT columns. `CollectV1VariantPathFilesFromLegacyRowsetSchema` pins this: a schema with the parent index and two path suffixes and no `index_info` reports all three files and their exact byte total. If there is a writer that leaves a suffixed V1 file that neither `index_info` nor the rowset schema lists, please point me to it and I will add a segment-derived enumeration. ########## be/src/format/table/index_disk_usage_reader.cpp: ########## @@ -0,0 +1,373 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include "format/table/index_disk_usage_reader.h" + +#include <boost/algorithm/string/case_conv.hpp> +#include <shared_mutex> +#include <string> +#include <string_view> +#include <utility> +#include <variant> + +#include "cloud/cloud_tablet.h" +#include "common/cast_set.h" +#include "core/block/block.h" +#include "core/column/column_nullable.h" +#include "core/column/column_string.h" +#include "core/column/column_vector.h" +#include "runtime/exec_env.h" +#include "runtime/runtime_state.h" +#include "storage/rowset/rowset.h" +#include "storage/tablet/base_tablet.h" +#include "storage/tablet/tablet_schema.h" + +namespace doris { + +using segment_v2::IndexDiskUsageLevel; +using segment_v2::IndexDiskUsageRecord; +using segment_v2::IndexDiskUsageRow; +using segment_v2::IndexDiskUsageStructure; + +namespace { + +const std::vector<std::pair<std::string_view, IndexDiskUsageReader::Column>>& column_names() { + using C = IndexDiskUsageReader::Column; + static const std::vector<std::pair<std::string_view, C>> names = { + {"PARTITION_NAME", C::kPartitionName}, + {"MATERIALIZED_INDEX_NAME", C::kMaterializedIndexName}, + {"TABLET_ID", C::kTabletId}, + {"BACKEND_ID", C::kBackendId}, + {"ROWSET_ID", C::kRowsetId}, + {"SEGMENT_ID", C::kSegmentId}, + {"INDEX_ID", C::kIndexId}, + {"INDEX_NAME", C::kIndexName}, + {"INDEX_TYPE", C::kIndexType}, + {"COLUMN_NAME", C::kColumnName}, + {"INDEX_SUFFIX", C::kIndexSuffix}, + {"STRUCTURE", C::kStructure}, + {"STORAGE_FORMAT", C::kStorageFormat}, + {"SEGMENT_COUNT", C::kSegmentCount}, + {"ROW_COUNT", C::kRowCount}, + {"TOTAL_BYTES", C::kTotalBytes}, + {"DICT_BYTES", C::kDictBytes}, + {"POSTING_BYTES", C::kPostingBytes}, + {"POSITION_BYTES", C::kPositionBytes}, + {"STATS_BYTES", C::kStatsBytes}, + {"OTHER_BYTES", C::kOtherBytes}, + {"STATS_SOURCE", C::kStatsSource}, + }; + return names; +} + +Result<IndexDiskUsageReader::Column> column_of(const std::string& slot_name) { + const std::string upper = boost::to_upper_copy(slot_name); + for (const auto& [name, column] : column_names()) { + if (upper == name) { + return column; + } + } + return ResultError(Status::InternalError("unknown index_disk_usage column {}", slot_name)); +} + +Result<IndexDiskUsageLevel> parse_level(const std::string& level) { + if (level == "tablet") { + return IndexDiskUsageLevel::kTablet; + } + if (level == "rowset") { + return IndexDiskUsageLevel::kRowset; + } + if (level == "segment") { + return IndexDiskUsageLevel::kSegment; + } + return ResultError(Status::InvalidArgument("unsupported index_disk_usage level {}", level)); +} + +std::string_view structure_name(IndexDiskUsageStructure structure) { + switch (structure) { + case IndexDiskUsageStructure::kTerm: + return "TERM"; + case IndexDiskUsageStructure::kBkd: + return "BKD"; + case IndexDiskUsageStructure::kAnn: + return "ANN"; + case IndexDiskUsageStructure::kContainer: + return "CONTAINER"; + } + return "UNKNOWN"; +} + +void insert_null(IColumn* column) { + auto& nullable = reinterpret_cast<ColumnNullable&>(*column); + nullable.get_nested_column().insert_default(); + nullable.get_null_map_data().push_back(1); +} + +IColumn* non_null_nested(IColumn* column) { + auto& nullable = reinterpret_cast<ColumnNullable&>(*column); + nullable.get_null_map_data().push_back(0); + return nullable.get_nested_column_ptr().get(); +} + +void insert_int64(IColumn* column, int64_t value) { + assert_cast<ColumnInt64*>(non_null_nested(column))->insert_value(value); +} + +void insert_int32(IColumn* column, int32_t value) { + assert_cast<ColumnInt32*>(non_null_nested(column))->insert_value(value); +} + +void insert_string(IColumn* column, std::string_view value) { + assert_cast<ColumnString*>(non_null_nested(column))->insert_data(value.data(), value.size()); +} + +} // namespace + +IndexDiskUsageReader::IndexDiskUsageReader(std::vector<SlotDescriptor*> slots, RuntimeState* state, + RuntimeProfile* /*profile*/, TMetaScanRange scan_range) + : _state(state), _slots(std::move(slots)), _scan_range(std::move(scan_range)) {} + +Status IndexDiskUsageReader::init_reader() { + if (!_scan_range.__isset.index_disk_usage_params) { + return Status::InvalidArgument("index_disk_usage scan range has no parameters"); + } + const TIndexDiskUsageMetadataParams& params = _scan_range.index_disk_usage_params; + _level = DORIS_TRY(parse_level(params.level)); + _options.position_detail = params.position_detail; + _options.index_ids.insert(params.index_ids.begin(), params.index_ids.end()); + _options.check_cancelled = [state = _state]() { + RETURN_IF_CANCELLED(state); + return Status::OK(); + }; + _slot_columns.clear(); + for (const SlotDescriptor* slot : _slots) { + const Column column = DORIS_TRY(column_of(slot->col_name())); + _slot_columns.push_back(column); + } + return Status::OK(); +} + +Status IndexDiskUsageReader::_do_get_next_block(Block* block, size_t* read_rows, bool* eof) { + const auto& tablets = _scan_range.index_disk_usage_params.tablets; + *read_rows = 0; + while (_next_tablet < tablets.size()) { + RETURN_IF_CANCELLED(_state); + const TIndexDiskUsageTablet& target = tablets[_next_tablet++]; + std::vector<IndexDiskUsageRow> rows; + TabletSchemaSPtr current_schema; + RETURN_IF_ERROR(_collect_tablet(target, &rows, ¤t_schema)); + rows = segment_v2::aggregate_index_disk_usage(std::move(rows), _level); + if (rows.empty()) { + continue; + } + RETURN_IF_ERROR(_fill_block(block, target, *current_schema, rows)); + *read_rows = rows.size(); + *eof = false; + return Status::OK(); + } + *eof = true; + return Status::OK(); +} + +Status IndexDiskUsageReader::_collect_tablet(const TIndexDiskUsageTablet& target, + std::vector<IndexDiskUsageRow>* rows, + TabletSchemaSPtr* current_schema) const { + BaseTabletSPtr tablet = DORIS_TRY(ExecEnv::get_tablet(target.tablet_id)); + if (auto cloud_tablet = std::dynamic_pointer_cast<CloudTablet>(tablet)) { + SyncOptions options; + options.query_version = target.version; + RETURN_IF_ERROR(cloud_tablet->sync_rowsets(options)); Review Comment: Fixed in 4f2458f7423. My earlier "Not changed" reply was wrong about the risk. The scan now calls `CloudTablet::sync_rowsets()` without a query version, the same call `calc_file_crc_action` uses, so it no longer stops at a cached `_max_version` and the cached compaction counters reach meta-service. A same-version compaction done on another backend is therefore picked up. The concern that this would make `capture_consistent_rowsets` fail when a compaction crosses the target version does not hold: `CloudTablet::add_rowsets` with version overlap moves the replaced rowsets to `_stale_rs_version_map`, and `capture_consistent_rowsets_unlocked` falls back to them (`include_stale_rowsets` defaults to true). Delete bitmaps are still synchronized, so a later MoW query on the same backend is unaffected. The cost is one meta-service call per tablet per scan, which is acceptable for a diagnostic function. `IndexDiskUsageCaptureTest.CaptureSyncsCompactionAtCachedVersion` returned the pre-compaction rowsets `[2-2],[3-3]` without syncing before the change and returns `[2-3]` now. `CaptureSurvivesCompactionPastScanVersion` shows that a compaction output `[2-4]` still lets version 3 be captured from the stale rowsets. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
