mrhhsg commented on code in PR #68032:
URL: https://github.com/apache/doris/pull/68032#discussion_r4119380679
##########
fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceProxy.java:
##########
@@ -97,6 +97,32 @@ public boolean needReconn() {
}
}
+ public Cloud.GetSpillStatsResponse
getSpillStats(Cloud.GetSpillStatsRequest request)
Review Comment:
Obsolete: `MetaServiceProxy.getSpillStats` was removed together with the
meta-service spill stats; the FE polls the BEs through `get_be_resource`
instead.
##########
be/src/common/config.cpp:
##########
@@ -1610,6 +1610,13 @@ DEFINE_String(spill_storage_limit, "20%");
// 20%
DEFINE_mInt32(spill_gc_interval_ms, "2000"); // 2s
DEFINE_mInt32(spill_gc_work_time_ms, "2000"); // 2s
DEFINE_mInt64(spill_file_part_size_bytes, "1073741824"); // 1GB
+DEFINE_String(spill_storage_type, "local");
+DEFINE_Validator(spill_storage_type, [](const std::string& config) -> bool {
+ return config == "local" || config == "s3";
+});
+DEFINE_String(spill_s3_storage_vault, "");
+DEFINE_mInt64(spill_s3_storage_limit_bytes, "0");
Review Comment:
Fixed: `spill_s3_storage_limit_bytes` has a validator requiring `>= 0`
(be/src/common/config.cpp).
##########
be/src/exec/spill/spill_file_manager.cpp:
##########
@@ -330,6 +571,61 @@ SpillDataDir::SpillDataDir(std::string path, int64_t
capacity_bytes,
INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity,
spill_disk_has_spill_gc_data);
}
+Status SpillDataDir::ensure_ready() {
+ if (!_is_remote || ready()) {
+ return Status::OK();
+ }
+ std::lock_guard<std::mutex> lock(_init_mutex);
+ if (ready()) {
+ return Status::OK();
+ }
+ if (!config::is_cloud_mode()) {
+ return Status::InternalError("spill to s3 is only supported in cloud
mode");
+ }
+ if (config::cloud_unique_id.empty()) {
+ return Status::InternalError(
+ "spill to s3 is not ready: cloud_unique_id is empty, waiting
for FE heartbeat");
+ }
+ // Resolve from what is already known locally; never trigger a
meta-service sync here.
+ // The vault refresh thread and the heartbeat fill these in, and callers
retry.
+ auto& engine = ExecEnv::GetInstance()->storage_engine().to_cloud();
+ std::string vault_id = _vault_id.empty() ? engine.default_vault_id() :
_vault_id;
+ io::RemoteFileSystemSPtr fs =
+ vault_id.empty() ? engine.latest_fs() :
doris::get_filesystem(vault_id);
+ if (fs == nullptr) {
+ return Status::InternalError(
+ "spill to s3 is not ready: storage vault '{}' not found (empty
means the default "
+ "vault of the instance; set spill_s3_storage_vault to the
vault ID in be.conf if "
+ "the instance has no default vault)",
+ vault_id);
+ }
+ if (fs->type() != io::FileSystemType::S3) {
+ return Status::NotSupported("spill to s3 only supports S3 storage
vaults, vault '{}' is {}",
+ vault_id, fs->type());
+ }
+ init_remote_fs(fs, config::cloud_unique_id);
+ return Status::OK();
+}
+
+void SpillDataDir::init_remote_fs(io::FileSystemSPtr fs, const std::string&
cloud_unique_id) {
+ DCHECK(_is_remote);
+ _fs = std::move(fs);
+ _remote_be_root = fmt::format("{}/{}", SPILL_DIR_PREFIX, cloud_unique_id);
Review Comment:
Obsolete: the object root no longer uses `cloud_unique_id` (it is
`spill/{ip}_{heartbeat_port}/{query_id}/...`), and the per-boot cleanup that
classified other boots as old was removed in 422b839.
--
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]