github-actions[bot] commented on code in PR #67977: URL: https://github.com/apache/doris/pull/67977#discussion_r4214577983
########## 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, ¤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)); + 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: [P2] Avoid a second cloud rowset synchronization on a cold tablet lookup. `ExecEnv::get_tablet()` enters `CloudTabletMgr::get_tablet`, whose cache-miss loader has already called `CloudMetaMgr::sync_tablet_rowsets` before it returns. This unversioned `sync_rowsets()` bypasses the query-version early exit and sends another synchronous `get_rowset` RPC for the same tablet even when nothing has changed. A cold scan can cover up to 20,000 tablets, so the repeated RPCs add material latency and meta-service load. Track whether the lookup just synchronized the tablet and skip this refresh on that path, while keeping the cache-hit compaction refresh; cover both paths with an RPC-count test. -- 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]
