mrhhsg commented on code in PR #68032: URL: https://github.com/apache/doris/pull/68032#discussion_r4228872272
########## be/src/exec/spill/remote_spill_data_dir.cpp: ########## @@ -0,0 +1,133 @@ +// 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 "exec/spill/remote_spill_data_dir.h" + +#include <glog/logging.h> + +#include <utility> + +#include "cloud/cloud_storage_engine.h" +#include "cloud/config.h" +#include "common/config.h" +#include "common/logging.h" +#include "common/metrics/metrics.h" +#include "io/fs/remote_file_system.h" +#include "runtime/exec_env.h" +#include "service/backend_options.h" +#include "storage/olap_define.h" +#include "storage/storage_policy.h" +#include "util/pretty_printer.h" + +namespace doris { + +RemoteSpillDataDir::RemoteSpillDataDir(std::string vault_id) + : SpillDataDir(fmt::format("s3:{}", vault_id.empty() ? "default" : vault_id), + /*spill_root=*/"", + fmt::format("s3:{}", vault_id.empty() ? "default" : vault_id), + /*capacity_bytes=*/0, TStorageMedium::S3), + _vault_id(std::move(vault_id)) {} + +Status RemoteSpillDataDir::init() { + RETURN_IF_ERROR(update_capacity()); + LOG(INFO) << fmt::format("remote spill store registered, vault_id={}, limit={}", + _vault_id.empty() ? "<default>" : _vault_id, + PrettyPrinter::print_bytes(_spill_data_limit_bytes)); + return Status::OK(); +} + +Status RemoteSpillDataDir::ensure_ready() { + if (ready()) { Review Comment: 已在 d1b5f50489c 修复:创建每个新 spill 文件时重新解析当前默认 vault,并将文件系统引用固定到文件、writer/reader 和失败删除队列。新增 A→B 轮换测试,覆盖两个文件同时存活、旧文件读取以及跨 vault 的 GC 重试;BE ASAN 定向测试 72/72 通过。 ########## be/src/exec/spill/spill_file_manager.cpp: ########## @@ -308,145 +456,8 @@ void SpillFileManager::gc(int32_t max_work_time_ms) { } } -DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(spill_disk_capacity, MetricUnit::BYTES); -DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(spill_disk_limit, MetricUnit::BYTES); -DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(spill_disk_avail_capacity, MetricUnit::BYTES); -DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(spill_disk_data_size, MetricUnit::BYTES); -DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(spill_disk_has_spill_data, MetricUnit::BYTES); -DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(spill_disk_has_spill_gc_data, MetricUnit::BYTES); - -SpillDataDir::SpillDataDir(std::string path, int64_t capacity_bytes, - TStorageMedium::type storage_medium) - : _path(std::move(path)), - _disk_capacity_bytes(capacity_bytes), - _storage_medium(storage_medium) { - spill_data_dir_metric_entity = DorisMetrics::instance()->metric_registry()->register_entity( - std::string("spill_data_dir.") + _path, {{"path", _path + "/" + SPILL_DIR_PREFIX}}); - INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity, spill_disk_capacity); - INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity, spill_disk_limit); - INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity, spill_disk_avail_capacity); - INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity, spill_disk_data_size); - INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity, spill_disk_has_spill_data); - INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity, spill_disk_has_spill_gc_data); -} - -bool is_directory_empty(const std::filesystem::path& dir) { - // Spill cleanup may delete the directory while the iterator is constructed or advanced. Treat - // that race as empty for these presence metrics. - try { - return std::filesystem::is_directory(dir) && - std::filesystem::directory_iterator(dir) == - std::filesystem::end(std::filesystem::directory_iterator {}); - } catch (const std::filesystem::filesystem_error&) { - return true; - } -} - -Status SpillDataDir::init() { - bool exists = false; - RETURN_IF_ERROR(io::global_local_filesystem()->exists(_path, &exists)); - if (!exists) { - RETURN_NOT_OK_STATUS_WITH_WARN(Status::IOError("opendir failed, path={}", _path), - "check file exist failed"); - } - RETURN_IF_ERROR(update_capacity()); - LOG(INFO) << fmt::format( - "spill storage path: {}, capacity: {}, limit: {}, available: " - "{}", - _path, PrettyPrinter::print_bytes(_disk_capacity_bytes), - PrettyPrinter::print_bytes(_spill_data_limit_bytes), - PrettyPrinter::print_bytes(_available_bytes)); - return Status::OK(); -} - -std::string SpillDataDir::get_spill_data_path(const std::string& query_id) const { - auto dir = fmt::format("{}/{}", _path, SPILL_DIR_PREFIX); - if (!query_id.empty()) { - dir = fmt::format("{}/{}", dir, query_id); - } - return dir; -} - -std::string SpillDataDir::get_spill_data_gc_path(const std::string& sub_dir_name) const { - auto dir = fmt::format("{}/{}", _path, SPILL_GC_DIR_PREFIX); - if (!sub_dir_name.empty()) { - dir = fmt::format("{}/{}", dir, sub_dir_name); - } - return dir; +int64_t SpillFileManager::remote_spill_data_bytes() { Review Comment: 已在 d1b5f50489c 将容量预留与已完成对象字节分开:容量限制仍保守预留,SHOW DATA 只取成功关闭的 part 字节;文件删除或 GC 重试成功后再扣减。新增 writer 尚未 close 时预留量大于零、已存对象量为零的测试,以及失败上传/删除路径断言。BE ASAN 定向测试 72/72 通过。 -- 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]
