mrhhsg commented on code in PR #68032:
URL: https://github.com/apache/doris/pull/68032#discussion_r4119377944
##########
be/src/cloud/cloud_meta_mgr.cpp:
##########
@@ -1893,7 +1900,30 @@ Status CloudMetaMgr::finish_restore_job(const int64_t
tablet_id, bool is_complet
});
}
-Status CloudMetaMgr::get_storage_vault_info(StorageVaultInfos* vault_infos,
bool* is_vault_mode) {
+Status CloudMetaMgr::report_spill_stats(int64_t boot_id, int64_t
remote_write_bytes,
+ int64_t remote_put_requests) {
+ ReportSpillStatsRequest req;
+ ReportSpillStatsResponse resp;
+ req.set_cloud_unique_id(config::cloud_unique_id);
+ auto* stats = req.mutable_stats();
+ stats->set_cloud_unique_id(config::cloud_unique_id);
Review Comment:
Obsolete: `report_spill_stats` and the reporter identity were removed;
nothing is reported to the meta-service any more.
##########
cloud/src/meta-service/meta_service.cpp:
##########
@@ -3530,6 +3530,143 @@ void
MetaServiceImpl::get_tablet_stats(::google::protobuf::RpcController* contro
}
}
+void MetaServiceImpl::report_spill_stats(::google::protobuf::RpcController*
controller,
+ const ReportSpillStatsRequest*
request,
+ ReportSpillStatsResponse* response,
+ ::google::protobuf::Closure* done) {
+ RPC_PREPROCESS(report_spill_stats, put);
+ instance_id = get_instance_id(resource_mgr_, request->cloud_unique_id());
+ if (instance_id.empty()) {
+ code = MetaServiceCode::INVALID_ARGUMENT;
+ msg = "empty instance_id";
+ LOG(INFO) << msg << ", cloud_unique_id=" << request->cloud_unique_id();
+ return;
+ }
+ RPC_RATE_LIMIT(report_spill_stats)
+ if (!request->has_stats() || request->stats().cloud_unique_id().empty()) {
+ code = MetaServiceCode::INVALID_ARGUMENT;
+ msg = "spill stats without cloud_unique_id";
+ return;
+ }
+ const auto& spill_stats = request->stats();
+ if (spill_stats.boot_id() <= 0 || spill_stats.remote_write_bytes() < 0 ||
+ spill_stats.remote_put_requests() < 0) {
+ code = MetaServiceCode::INVALID_ARGUMENT;
+ msg = fmt::format("invalid spill stats, boot_id={} write_bytes={}
put_requests={}",
+ spill_stats.boot_id(),
spill_stats.remote_write_bytes(),
+ spill_stats.remote_put_requests());
+ return;
+ }
+
+ TxnErrorCode err = txn_kv_->create_txn(&txn);
+ if (err != TxnErrorCode::TXN_OK) {
+ code = cast_as<ErrCategory::CREATE>(err);
+ msg = fmt::format("failed to create txn, err={}", err);
+ return;
+ }
+ // One record per BE. The report carries totals since boot, so a report of
the same boot
+ // replaces the previous one (a retry cannot double count). The first
report of a new boot
+ // folds the previous process' totals into prior_boots_*, which keeps the
record count
+ // bounded by the number of BEs instead of the number of BE restarts.
+ std::string key = stats_spill_key({instance_id,
spill_stats.cloud_unique_id()});
+ std::string existing_val;
+ err = txn->get(key, &existing_val);
+ if (err != TxnErrorCode::TXN_OK && err != TxnErrorCode::TXN_KEY_NOT_FOUND)
{
+ code = cast_as<ErrCategory::READ>(err);
+ msg = fmt::format("failed to get spill stats, err={}", err);
+ return;
+ }
+ SpillStatsPB value = spill_stats;
+ if (err == TxnErrorCode::TXN_OK) {
+ SpillStatsPB existing;
+ if (!existing.ParseFromString(existing_val)) {
+ code = MetaServiceCode::PROTOBUF_PARSE_ERR;
+ msg = fmt::format("malformed spill stats, key={}", hex(key));
+ return;
+ }
+ int64_t prior_bytes = existing.prior_boots_write_bytes();
+ int64_t prior_requests = existing.prior_boots_put_requests();
+ if (existing.boot_id() != spill_stats.boot_id()) {
+ prior_bytes += existing.remote_write_bytes();
+ prior_requests += existing.remote_put_requests();
+ } else {
+ // Totals of one process only grow; a late-delivered duplicate
must not roll back.
+ value.set_remote_write_bytes(
+ std::max(existing.remote_write_bytes(),
spill_stats.remote_write_bytes()));
+ value.set_remote_put_requests(
+ std::max(existing.remote_put_requests(),
spill_stats.remote_put_requests()));
+ }
+ value.set_prior_boots_write_bytes(prior_bytes);
+ value.set_prior_boots_put_requests(prior_requests);
+ }
+
value.set_update_time_ms(std::chrono::duration_cast<std::chrono::milliseconds>(
+
std::chrono::system_clock::now().time_since_epoch())
+ .count());
+ txn->put(key, value.SerializeAsString());
+ err = txn->commit();
+ if (err != TxnErrorCode::TXN_OK) {
+ code = cast_as<ErrCategory::COMMIT>(err);
+ msg = fmt::format("failed to commit spill stats, err={}
cloud_unique_id={} boot_id={}", err,
+ spill_stats.cloud_unique_id(),
spill_stats.boot_id());
+ return;
+ }
+}
+
+void MetaServiceImpl::get_spill_stats(::google::protobuf::RpcController*
controller,
+ const GetSpillStatsRequest* request,
+ GetSpillStatsResponse* response,
+ ::google::protobuf::Closure* done) {
+ RPC_PREPROCESS(get_spill_stats, get);
+ instance_id = get_instance_id(resource_mgr_, request->cloud_unique_id());
+ if (instance_id.empty()) {
+ code = MetaServiceCode::INVALID_ARGUMENT;
+ msg = "empty instance_id";
+ LOG(INFO) << msg << ", cloud_unique_id=" << request->cloud_unique_id();
+ return;
+ }
+ RPC_RATE_LIMIT(get_spill_stats)
+
+ TxnErrorCode err = txn_kv_->create_txn(&txn);
+ if (err != TxnErrorCode::TXN_OK) {
+ code = cast_as<ErrCategory::CREATE>(err);
+ msg = fmt::format("failed to create txn, err={}", err);
+ return;
+ }
+ std::string key0 = stats_spill_key_prefix(instance_id);
Review Comment:
Obsolete: the per-identity spill stats keys in the meta-service were removed.
##########
be/src/exec/spill/spill_file_manager.cpp:
##########
@@ -123,6 +213,24 @@ void SpillFileManager::_spill_gc_thread_callback() {
for (auto& [path, dir] : _spill_store_map) {
static_cast<void>(dir->update_capacity());
}
+ // Both configs are mutable; observe changes without a restart.
+ int64_t budget_limit = config::spill_s3_max_inflight_upload_bytes;
Review Comment:
Fixed earlier: startup no longer treats `limit < 2 * s3_write_buffer_size`
as fatal. Startup and the dynamic update path now apply the same rule: warn and
serialize uploads (see doris_main.cpp and `_spill_gc_thread_callback`).
--
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]