airborne12 commented on code in PR #67977:
URL: https://github.com/apache/doris/pull/67977#discussion_r4216352857


##########
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:
   Fixed in bbde70de110. My earlier reply argued that the rowset schema already 
names these files. I could not find a writer that breaks that, but the physical 
files are the ground truth, so index_disk_usage now looks for them directly.
   
   For a V1 rowset without `index_info` whose schema has a VARIANT column and 
that lives on a local filesystem, the segment directory is searched for 
`{rowset_id}_{segment_id}_{index_id}[@suffix].idx` files that are not listed 
yet, in index id and suffix order. Rowsets with `index_info`, without a VARIANT 
column, or on a non-local filesystem never list the directory. The tablet's 
segments share one sorted, names-only listing (no per-file stat, since unused 
rowsets delete files from the same directory while it is read), and each 
segment finds its files by binary search.
   
   `CollectV1VariantPathFilesFromDirectory` uses a schema that lists only the 
parent index, no `index_info`, and files for the parent, two paths and a path 
with a two-digit index id, next to look-alike files: another segment, a segment 
whose id shares the prefix, a V2 container name, a non-numeric and a negative 
id, and another extension. It reports exactly the four files with their byte 
total, and it found only the parent when the fallback was disabled. 
`CollectV1WithoutVariantColumnDoesNotListDirectory`, 
`CollectV1SegmentsShareOneDirectoryListing` and the `DirectoryFileNames` tests 
cover the gating, the sharing and the listing itself.
   



##########
fe/fe-core/src/main/java/org/apache/doris/tablefunction/IndexDiskUsageTableValuedFunction.java:
##########
@@ -0,0 +1,387 @@
+// 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.
+
+package org.apache.doris.tablefunction;
+
+import org.apache.doris.analysis.TupleDescriptor;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Database;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.Index;
+import org.apache.doris.catalog.MaterializedIndex;
+import org.apache.doris.catalog.MaterializedIndex.IndexExtState;
+import org.apache.doris.catalog.OlapTable;
+import org.apache.doris.catalog.Partition;
+import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.TableIf;
+import org.apache.doris.catalog.Tablet;
+import org.apache.doris.catalog.info.IndexType;
+import org.apache.doris.cloud.catalog.CloudPartition;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.ErrorCode;
+import org.apache.doris.common.Pair;
+import org.apache.doris.datasource.InternalCatalog;
+import org.apache.doris.datasource.tvf.source.IndexDiskUsageScanNode;
+import org.apache.doris.mysql.privilege.PrivPredicate;
+import org.apache.doris.nereids.exceptions.AnalysisException;
+import org.apache.doris.planner.PlanNodeId;
+import org.apache.doris.planner.ScanContext;
+import org.apache.doris.planner.ScanNode;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.SessionVariable;
+import org.apache.doris.rpc.RpcException;
+import org.apache.doris.thrift.TIndexDiskUsageMetadataParams;
+import org.apache.doris.thrift.TIndexDiskUsageTablet;
+import org.apache.doris.thrift.TMetaScanRange;
+import org.apache.doris.thrift.TMetadataType;
+
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableSet;
+import com.google.common.collect.Lists;
+import com.google.common.collect.Maps;
+import org.apache.commons.lang3.StringUtils;
+
+import java.util.Arrays;
+import java.util.Collection;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+/**
+ * The implement of table valued function
+ * index_disk_usage("database" = "db1", "table" = "table1").
+ * It reports the physical bytes of every inverted index of the table, split 
by component.
+ */
+public class IndexDiskUsageTableValuedFunction extends 
MetadataTableValuedFunction {
+    public static final String NAME = "index_disk_usage";
+
+    private static final String DATABASE = "database";
+    private static final String TABLE = "table";
+    private static final String PARTITIONS = "partitions";
+    private static final String INDEXES = "indexes";
+    private static final String LEVEL = "level";
+    private static final String POSITION_DETAIL = "position_detail";
+
+    private static final ImmutableSet<String> PROPERTIES_SET =
+            ImmutableSet.of(DATABASE, TABLE, PARTITIONS, INDEXES, LEVEL, 
POSITION_DETAIL);
+    private static final ImmutableSet<String> LEVELS = 
ImmutableSet.of("tablet", "rowset", "segment");
+
+    private static final ImmutableList<Column> SCHEMA = ImmutableList.of(
+            varcharColumn("PARTITION_NAME"),
+            varcharColumn("MATERIALIZED_INDEX_NAME"),
+            bigintColumn("TABLET_ID"),
+            bigintColumn("BACKEND_ID"),
+            varcharColumn("ROWSET_ID"),
+            new Column("SEGMENT_ID", ScalarType.createType(PrimitiveType.INT), 
true),
+            bigintColumn("INDEX_ID"),
+            varcharColumn("INDEX_NAME"),
+            varcharColumn("INDEX_TYPE"),
+            varcharColumn("COLUMN_NAME"),
+            varcharColumn("INDEX_SUFFIX"),
+            varcharColumn("STRUCTURE"),
+            varcharColumn("STORAGE_FORMAT"),
+            bigintColumn("SEGMENT_COUNT"),
+            bigintColumn("ROW_COUNT"),
+            bigintColumn("TOTAL_BYTES"),
+            bigintColumn("DICT_BYTES"),
+            bigintColumn("POSTING_BYTES"),
+            bigintColumn("POSITION_BYTES"),
+            bigintColumn("STATS_BYTES"),
+            bigintColumn("OTHER_BYTES"),
+            varcharColumn("STATS_SOURCE"));
+
+    /**
+     * A tablet of a base or rollup index to inspect, pinned to the visible 
version of its partition.
+     */
+    public static class TabletTarget {
+        private final Tablet tablet;
+        private final long partitionId;
+        private final long materializedIndexId;
+        private final long version;
+
+        public TabletTarget(Tablet tablet, long partitionId, long 
materializedIndexId, long version) {
+            this.tablet = tablet;
+            this.partitionId = partitionId;
+            this.materializedIndexId = materializedIndexId;
+            this.version = version;
+        }
+
+        public long getMaterializedIndexId() {
+            return materializedIndexId;
+        }
+
+        public Tablet getTablet() {
+            return tablet;
+        }
+
+        public long getTabletId() {
+            return tablet.getId();
+        }
+
+        public long getPartitionId() {
+            return partitionId;
+        }
+
+        public long getVersion() {
+            return version;
+        }
+
+        public TIndexDiskUsageTablet toThrift() {
+            TIndexDiskUsageTablet target = new TIndexDiskUsageTablet();
+            target.setTabletId(getTabletId());
+            target.setPartitionId(partitionId);
+            target.setMaterializedIndexId(materializedIndexId);
+            target.setVersion(version);
+            return target;
+        }
+    }
+
+    private final String level;
+    private final boolean positionDetail;
+    private final List<Long> indexIds;
+    private final Map<Long, String> partitionNames;
+    private final Map<Long, String> materializedIndexNames;
+    private final List<TabletTarget> tabletTargets;
+
+    public IndexDiskUsageTableValuedFunction(Map<String, String> params) 
throws AnalysisException {
+        Map<String, String> validParams = Maps.newHashMap();
+        for (Map.Entry<String, String> entry : params.entrySet()) {
+            String key = entry.getKey().toLowerCase();
+            if (!PROPERTIES_SET.contains(key)) {
+                throw new AnalysisException("'" + entry.getKey() + "' is 
invalid property");
+            }
+            validParams.put(key, entry.getValue());
+        }
+        String dbName = validParams.get(DATABASE);
+        String tableName = validParams.get(TABLE);
+        if (StringUtils.isEmpty(dbName) || StringUtils.isEmpty(tableName)) {
+            throw new AnalysisException("'database' and 'table' are required 
for index_disk_usage");
+        }
+        this.level = parseLevel(validParams.getOrDefault(LEVEL, "tablet"));
+        this.positionDetail = 
parsePositionDetail(validParams.getOrDefault(POSITION_DETAIL, "false"));
+        checkShowPrivilege(dbName, tableName);
+
+        OlapTable table = getOlapTable(dbName, tableName);
+        String qualifiedName = dbName + "." + tableName;
+        List<Long> resolvedIndexIds;
+        List<Partition> partitions;
+        Map<Long, String> resolvedPartitionNames = Maps.newLinkedHashMap();
+        Map<Long, String> resolvedMaterializedIndexNames = 
Maps.newLinkedHashMap();
+        Map<Long, List<Pair<Long, List<Tablet>>>> tabletsByPartition = 
Maps.newHashMap();
+        table.readLock();
+        try {
+            resolvedIndexIds = resolveIndexIds(table, 
validParams.get(INDEXES), qualifiedName);
+            partitions = resolvePartitions(table, validParams.get(PARTITIONS), 
qualifiedName);
+            for (Partition partition : partitions) {
+                resolvedPartitionNames.put(partition.getId(), 
partition.getName());
+                // A light ADD INDEX also installs indexes on rollups, so 
their tablets can hold index files.
+                List<Pair<Long, List<Tablet>>> indexTablets = 
Lists.newArrayList();
+                for (MaterializedIndex index : 
partition.getMaterializedIndices(IndexExtState.VISIBLE)) {
+                    resolvedMaterializedIndexNames.putIfAbsent(index.getId(), 
table.getIndexNameById(index.getId()));
+                    indexTablets.add(Pair.of(index.getId(), 
Lists.newArrayList(index.getTablets())));
+                }
+                tabletsByPartition.put(partition.getId(), indexTablets);
+            }
+        } finally {
+            table.readUnlock();
+        }
+        // Cloud partitions may fetch visible versions from meta-service, so 
read them without the table lock.
+        List<Long> versions = visibleVersions(partitions);

Review Comment:
   Fixed in 4b594d4fce8. After the cloud versions are fetched, the FE takes the 
table read lock again and compares the selected partitions and the set of their 
visible materialized index ids with the copy. If a partition is gone or a 
rollup was dropped or became visible, the snapshot is taken again (3 attempts 
in total, then the query fails asking for a retry), so tablets are never paired 
with a version that was fetched for a different layout. The version source is a 
package-private constructor parameter so the cloud path can be unit tested.
   
   `testCloudSnapshotIsTakenAgainWhenRollupIsDroppedDuringVersionFetch` drops a 
rollup while the versions are read and expects two fetches and no tablet of the 
dropped rollup. `testCloudSnapshotFailsWhenRollupsKeepChanging` expects the 
bounded failure, `testCloudSnapshotReportsPartitionDroppedDuringVersionFetch` 
expects the unknown-partition error, and 
`testCloudSnapshotFetchesVersionsOnceWhenLayoutIsStable` expects a single 
fetch. With the re-check disabled the first three fail.
   



##########
be/src/format/table/index_disk_usage_reader.cpp:
##########
@@ -0,0 +1,417 @@
+// 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/query_context.h"
+#include "runtime/runtime_state.h"
+#include "storage/rowset/rowset.h"
+#include "storage/rowset/rowset_meta.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();
+    };
+    // Index file reads are accounted to the query like the reads of any other 
scan.
+    _io_ctx = io::IOContext {
+            .reader_type = ReaderType::READER_QUERY,
+            .query_id = &_state->query_id(),
+            .is_inverted_index = true,
+    };
+    _io_ctx.inverted_index_snii_read_no_write_file_cache =
+            
_state->query_options().inverted_index_snii_read_no_write_file_cache;
+    if (auto* query_ctx = _state->get_query_ctx(); query_ctx != nullptr) {
+        _io_ctx.remote_scan_cache_write_limiter = 
query_ctx->remote_scan_cache_write_limiter();
+    }
+    _options.io_ctx = &_io_ctx;
+    _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));
+    const std::vector<RowsetSharedPtr> rowsets = 
DORIS_TRY(capture_rowsets(tablet, target.version));
+    *current_schema = label_schema(tablet->tablet_schema(), rowsets);
+
+    const io::IOContext io_ctx = tablet_io_context(_io_ctx, 
tablet->ttl_seconds());
+    segment_v2::IndexDiskUsageOptions options = _options;
+    options.io_ctx = &io_ctx;
+    for (const RowsetSharedPtr& rowset : rowsets) {
+        RETURN_IF_ERROR(segment_v2::collect_rowset_index_disk_usage(rowset, 
options,
+                                                                    
target.tablet_id, rows));
+    }
+    return Status::OK();
+}
+
+Result<std::vector<RowsetSharedPtr>> IndexDiskUsageReader::capture_rowsets(
+        const BaseTabletSPtr& tablet, int64_t version) {
+    if (auto cloud_tablet = std::dynamic_pointer_cast<CloudTablet>(tablet)) {
+        // A compaction elsewhere may replace cached rowsets without changing 
the visible version,
+        // so this does not stop at a cached `version` the way a query sync 
does.
+        RETURN_IF_ERROR_RESULT(cloud_tablet->sync_rowsets());

Review Comment:
   Fixed in 3f1d4356672. `ExecEnv::get_tablet` now receives a 
`SyncRowsetStats`. A tablet that was already cached keeps the unversioned 
refresh, so a compaction done elsewhere at the cached version is still seen. A 
tablet that the lookup loaded only syncs up to the scan version, which sends no 
RPC when the load already saw that version. I kept that sync instead of 
skipping it, because a lookup can join a load that began before the scan 
version was committed (`tablet_meta_cache_miss` is counted before the 
single-flight load).
   
   `IndexDiskUsageCaptureTest` installs a `CloudStorageEngine` and counts 
`CloudMetaMgr::sync_tablet_rowsets` calls: `ColdLookupSyncsOnce` expects one 
sync for a cold lookup, `ColdLookupSyncsToTheScanVersionWhenTheLoadIsOlder` 
expects a second one only when the loaded tablet lacks the scan version, 
`CachedLookupSyncsCompactionAtCachedVersion` expects one for a cached lookup 
that sees a same-version compaction, and 
`CompactionPastScanVersionKeepsScanVersionReadable` covers a compaction output 
past the scan version. Syncing a loaded tablet unconditionally breaks the 
counts of the cold and cached tests, and skipping the sync for a loaded tablet 
breaks `ColdLookupSyncsToTheScanVersionWhenTheLoadIsOlder`.
   



-- 
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