mrhhsg commented on code in PR #68032:
URL: https://github.com/apache/doris/pull/68032#discussion_r4119375258
##########
be/src/io/fs/s3_file_writer.cpp:
##########
@@ -93,13 +100,75 @@ S3FileWriter::~S3FileWriter() {
s3_file_being_written << -1;
}
+void S3FileWriter::_record_request(std::atomic<int64_t> RemoteWriteStats::*
counter,
+ int64_t elapsed_ns, bool ok, int64_t
uploaded_bytes) {
+ if (_remote_write_stats == nullptr) {
+ return;
+ }
+ ((*_remote_write_stats).*counter)++;
Review Comment:
Documented instead of changed (bdd4fad): the value is no longer persisted as
a warehouse PUT total (the meta-service stats were removed); it only feeds the
profile and BE metrics. `RemoteWriteStats` now states that it counts logical
requests and that retries inside the object storage client are not included.
##########
be/src/exec/spill/spill_file.cpp:
##########
@@ -43,23 +43,26 @@ SpillFile::~SpillFile() {
}
void SpillFile::gc() {
- bool exists = false;
- auto status = io::global_local_filesystem()->exists(_spill_dir, &exists);
- if (status.ok() && exists) {
- // Delete spill directory directly instead of moving it to a GC
directory.
- // This simplifies cleanup and avoids retaining spill data under a GC
path.
- status = io::global_local_filesystem()->delete_directory(_spill_dir);
+ if (_dir_created) {
+ // Delete the spill directory (or object key prefix) directly instead
of moving it to a
+ // GC directory. No existence check: for object storage a "directory"
never exists as an
+ // object, while deleting a missing local directory or an empty prefix
is a no-op.
+ auto fs = _data_dir->fs();
+ Status status = fs != nullptr ? fs->delete_directory(_spill_dir)
+ : Status::InternalError("spill store {}
is not ready",
+
_data_dir->path());
DBUG_EXECUTE_IF("fault_inject::spill_file::gc", {
status = Status::Error<INTERNAL_ERROR>("fault_inject spill_file gc
failed");
});
if (!status.ok()) {
LOG_EVERY_T(WARNING, 1) << fmt::format("failed to delete spill
data, dir {}, error: {}",
_spill_dir,
status.to_string());
}
+ _dir_created = false;
}
// Decrease spill data usage even if per-file cleanup failed. QueryContext
teardown deletes the
// whole query spill directory and retains failures for later retries.
- _data_dir->update_spill_data_usage(-_total_written_bytes);
+ _data_dir->release(_total_written_bytes);
Review Comment:
Fixed earlier: a failed deletion keeps the bytes charged and hands them to
the manager's retry (`retry_spill_directory_deletion`), which releases them
only after the deletion succeeds; covered by
`SpillFileS3Test.FailedDeletionKeepsCapacityCharged`.
##########
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()) {
Review Comment:
Obsolete: the meta-service boot/generation folding was removed with the
meta-service spill stats.
--
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]