mrhhsg commented on code in PR #68032: URL: https://github.com/apache/doris/pull/68032#discussion_r4119376543
########## be/src/exec/spill/spill_remote_upload_budget.cpp: ########## @@ -0,0 +1,88 @@ +// 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/spill_remote_upload_budget.h" + +#include <glog/logging.h> + +#include <chrono> + +#include "util/stopwatch.hpp" + +namespace doris { + +Status SpillRemoteUploadBudget::acquire(int64_t bytes, const std::function<bool()>& is_cancelled, + int64_t* wait_ns) { + MonotonicStopWatch watch; + watch.start(); + std::unique_lock<std::mutex> lock(_mutex); + while (_inflight_bytes > 0 && _inflight_bytes + bytes > _limit_bytes) { + if (is_cancelled && is_cancelled()) { + return Status::Cancelled("query cancelled while waiting for spill upload budget"); + } + _cv.wait_for(lock, std::chrono::milliseconds(100)); + } + _inflight_bytes += bytes; + _total_acquired_bytes += bytes; + if (wait_ns != nullptr) { + *wait_ns = static_cast<int64_t>(watch.elapsed_time()); + } + return Status::OK(); +} + +void SpillRemoteUploadBudget::release(int64_t bytes) { + { + std::lock_guard<std::mutex> lock(_mutex); + _inflight_bytes -= bytes; + _total_released_bytes += bytes; + DCHECK_GE(_inflight_bytes, 0) << "spill upload budget released more than acquired"; Review Comment: Fixed: `release()` now uses an always-on `DORIS_CHECK_GE(_inflight_bytes, 0)` without clamping. ########## be/src/exec/spill/spill_file_writer.cpp: ########## @@ -90,44 +145,179 @@ Status SpillFileWriter::_close_current_part(const std::shared_ptr<SpillFile>& sp _part_meta.append((const char*)&_part_max_sub_block_size, sizeof(_part_max_sub_block_size)); _part_meta.append((const char*)&_part_written_blocks, sizeof(_part_written_blocks)); - { + int64_t meta_size = _part_meta.size(); + // The footer must always be written so that the part can be closed; account it + // without checking the capacity limit. + Status status = _data_dir->try_reserve(meta_size, /*force=*/true); + if (status.ok()) { SCOPED_TIMER(_write_file_timer); - RETURN_IF_ERROR(_file_writer->append(_part_meta)); + status = _file_writer->append(_part_meta); + if (!status.ok()) { + _data_dir->release(meta_size); + } } - int64_t meta_size = _part_meta.size(); - _part_written_bytes += meta_size; - COUNTER_UPDATE(_write_file_total_size, meta_size); - if (_resource_ctx) { - _resource_ctx->io_context()->update_spill_write_bytes_to_local_storage(meta_size); - } - if (_write_file_current_size) { - COUNTER_UPDATE(_write_file_current_size, meta_size); - } - _data_dir->update_spill_data_usage(meta_size); - ExecEnv::GetInstance()->spill_file_mgr()->update_spill_write_bytes(meta_size); - // Incrementally update SpillFile's accounting so gc() can always - // decrement the correct amount, even if close() is never called. - if (spill_file) { - spill_file->update_written_bytes(meta_size); + if (status.ok()) { + _part_written_bytes += meta_size; + COUNTER_UPDATE(_write_file_total_size, meta_size); + if (_resource_ctx) { + if (_data_dir->is_remote()) { + _resource_ctx->io_context()->update_spill_write_bytes_to_remote_storage(meta_size); + } else { + _resource_ctx->io_context()->update_spill_write_bytes_to_local_storage(meta_size); + } + } + if (_write_file_current_size) { + COUNTER_UPDATE(_write_file_current_size, meta_size); + } + ExecEnv::GetInstance()->spill_file_mgr()->update_spill_write_bytes(meta_size); + // Incrementally update SpillFile's accounting so gc() can always + // decrement the correct amount, even if close() is never called. + if (spill_file) { + spill_file->update_written_bytes(meta_size); + } } - RETURN_IF_ERROR(_file_writer->close()); - _file_writer.reset(); + // Issue a non-blocking close so that the upload of this part overlaps with the next + // one. The part is confirmed later by _reap_closing_parts(), which also reconciles + // the upload budget, so it is queued even when something above failed. + ClosingPart part; + part.path = _current_part_path; + part.part_index = _current_part_index; + part.part_bytes = _part_written_bytes; + part.ledger = std::move(_part_ledger); + part.stats = std::move(_part_stats); + part.close_status = status.ok() ? _file_writer->close(/*non_block=*/true) : status; + part.writer = std::move(_file_writer); + _closing_parts.emplace_back(std::move(part)); // Advance to next part ++_current_part_index; - if (spill_file) { - spill_file->increment_part_count(); - } _part_written_blocks = 0; _part_written_bytes = 0; _part_max_sub_block_size = 0; _part_meta.clear(); + return status; +} + +Status SpillFileWriter::_reap_closing_parts(bool block, + const std::shared_ptr<SpillFile>& spill_file) { + Status first_error; + while (!_closing_parts.empty()) { + auto& part = _closing_parts.front(); + Status st = part.close_status; + if (st.ok()) { + st = part.writer->try_finish_close(); + if (st.is<ErrorCode::NEED_SEND_AGAIN>()) { + if (!block) { + break; + } + st = part.writer->close(); + } else if (st.is<ErrorCode::NOT_IMPLEMENTED_ERROR>()) { + // Writers without an async close protocol (local files) finished the work + // in close(true); close() only flips the state. + st = part.writer->state() == io::FileWriter::State::CLOSED ? Status::OK() + : part.writer->close(); + } + } + st = _finish_part(part, spill_file, st); + _closing_parts.erase(_closing_parts.begin()); + if (!st.ok() && first_error.ok()) { + first_error = st; + if (!block) { + break; + } + } + } + return first_error; +} + +Status SpillFileWriter::_finish_part(ClosingPart& part, + const std::shared_ptr<SpillFile>& spill_file, + Status close_status) { + MultipartUploadId upload = _multipart_upload_id(part.writer.get()); + if (!close_status.ok() && part.writer != nullptr && + part.writer->state() != io::FileWriter::State::CLOSED) { + // The part never reached a final state (footer failed, or close(true) could not be + // issued). Destroying the writer waits for every in-flight upload, so the ledger below + // is complete afterwards. close() is not used for this: on a cancelled query it would + // be refused by the upload gate right away and drain nothing. + part.writer.reset(); + } + + // Budget: everything the gate took for this part but the upload callback never gave + // back (buffers that failed before their upload started) is released here. The writer is + // in its final state at this point, so every callback that will ever fire has fired. + if (_budget != nullptr && part.ledger != nullptr) { + int64_t remaining = part.ledger->acquired.load() - part.ledger->released.load(); + DCHECK_GE(remaining, 0) << "upload callback released more than acquired, part=" + << part.path; + if (remaining > 0) { + _budget->release(remaining); + } + if (_remote_upload_wait_timer != nullptr) { + COUNTER_UPDATE(_remote_upload_wait_timer, part.ledger->wait_ns.load()); + } + } + + if (part.stats != nullptr) { + int64_t requests = part.stats->total_requests(); + int64_t data_requests = part.stats->put_object_requests + part.stats->upload_part_requests; + int64_t uploaded_bytes = part.stats->uploaded_bytes; + if (_remote_write_requests != nullptr) { + COUNTER_UPDATE(_remote_write_requests, requests); + COUNTER_UPDATE(_remote_upload_part_requests, data_requests); + COUNTER_UPDATE(_remote_upload_bytes, uploaded_bytes); + COUNTER_UPDATE(_remote_upload_timer, part.stats->request_time_ns.load()); + } + if (_resource_ctx) { + _resource_ctx->io_context()->update_spill_remote_write_requests(requests); + } + ExecEnv::GetInstance()->spill_file_mgr()->update_spill_remote_write(uploaded_bytes, + data_requests); + } + + if (!close_status.ok()) { + LOG(WARNING) << "failed to close spill part " << part.path << ": " << close_status; + _abort_multipart_upload(upload); + return close_status; + } + if (spill_file) { + spill_file->add_part(part.part_bytes); + } return Status::OK(); } +SpillFileWriter::MultipartUploadId SpillFileWriter::_multipart_upload_id(io::FileWriter* writer) { + auto* s3_writer = dynamic_cast<io::S3FileWriter*>(writer); + if (s3_writer == nullptr || s3_writer->upload_id().empty()) { + return {}; + } + return {.path = s3_writer->path().native(), + .bucket = s3_writer->bucket(), + .key = s3_writer->key(), + .upload_id = s3_writer->upload_id()}; +} + +void SpillFileWriter::_abort_multipart_upload(const MultipartUploadId& upload) { + if (upload.upload_id.empty()) { + return; + } + auto s3_fs = std::dynamic_pointer_cast<io::S3FileSystem>(_data_dir->fs()); + if (s3_fs == nullptr) { + return; + } + auto client = s3_fs->client_holder()->get(); + if (client == nullptr) { + return; + } + auto resp = client->abort_multipart_upload( Review Comment: Won't fix in Doris, by design: residue of a BE that crashed (or of an abort that failed) is left to the bucket lifecycle rule, which must expire the keys under `spill/` and abort incomplete multipart uploads; this is documented at `spill_storage_type` in config.cpp. The best-effort abort only shortens how long the parts are billed. ########## be/src/exec/spill/spill_file_manager.cpp: ########## @@ -308,6 +438,99 @@ void SpillFileManager::gc(int32_t max_work_time_ms) { } } +void SpillFileManager::_remote_gc(SpillDataDir* store) { + if (!store->ready()) { + // Retry about once a minute at the default 2s GC interval. ensure_ready() only reads + // what the vault refresh thread and the FE heartbeat already brought in. + if (_remote_not_ready_rounds++ % 30 != 0) { + return; + } + auto st = store->ensure_ready(); + if (!st.ok()) { + LOG(WARNING) << "remote spill store is not ready yet: " << st; + return; + } + } + _report_remote_spill_stats(store); + if (!remote_startup_cleanup_pending()) { + return; + } + auto st = _remote_startup_cleanup(store); + if (st.ok()) { + _remote_startup_cleanup_pending.store(false, std::memory_order_release); + } else { + LOG_EVERY_T(WARNING, 60) << "failed to clean up spill objects of previous boots, will " + "retry: " + << st; + } +} + +void SpillFileManager::flush_remote_spill_stats() { + for (auto& [path, store] : _spill_store_map) { + if (store->is_remote()) { + _report_remote_spill_stats(store.get(), /*final_report=*/true); + } + } +} + +void SpillFileManager::_report_remote_spill_stats(SpillDataDir* store, bool final_report) { + std::lock_guard<std::mutex> lock(_remote_report_mutex); + // About once a minute at the default 2s GC interval; a final report skips the cadence. + if (!final_report && _remote_report_rounds++ % 30 != 0) { + return; + } + if (!store->ready() || !config::is_cloud_mode()) { + return; + } + int64_t write_bytes = remote_write_bytes_since_boot(); + int64_t put_requests = remote_put_requests_since_boot(); + if (write_bytes == _reported_remote_write_bytes && + put_requests == _reported_remote_put_requests) { + return; + } + auto st = ExecEnv::GetInstance()->storage_engine().to_cloud().meta_mgr().report_spill_stats( + store->boot_id(), write_bytes, put_requests); + if (!st.ok()) { + LOG_EVERY_T(WARNING, 60) << "failed to report spill stats to meta-service" + << (final_report ? "" : ", will retry") << ": " << st; + return; + } + _reported_remote_write_bytes = write_bytes; + _reported_remote_put_requests = put_requests; +} + +Status SpillFileManager::_remote_startup_cleanup(SpillDataDir* store) { + auto fs = store->fs(); + const auto& be_root = store->get_remote_be_root(); + const auto current_boot_id = std::to_string(store->boot_id()); + + MonotonicStopWatch watch; + watch.start(); + std::vector<io::FileInfo> files; + bool exists = false; + RETURN_IF_ERROR(fs->list(be_root, true, &files, &exists)); Review Comment: Obsolete: the startup cleanup that listed the BE root 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]
