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, &current_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]

Reply via email to