github-actions[bot] commented on code in PR #68032:
URL: https://github.com/apache/doris/pull/68032#discussion_r4229103814
##########
be/src/exec/spill/spill_file_manager.cpp:
##########
@@ -314,19 +418,20 @@ void
SpillFileManager::_release_external_spill_session(ExternalSpillSession* spi
}
const auto query_dir =
spill_session->_data_dir->get_spill_data_path(spill_session->_query_id);
- std::lock_guard directory_lock(_pending_query_spill_directories_mutex);
+ std::lock_guard directory_lock(_pending_spill_directories_mutex);
auto it = _external_spill_directory_leases.find(query_dir);
- DCHECK(it != _external_spill_directory_leases.end());
- if (it == _external_spill_directory_leases.end()) {
- return;
- }
- DCHECK_GT(it->second, 0);
+ DORIS_CHECK(it != _external_spill_directory_leases.end());
+ DORIS_CHECK_GT(it->second, 0);
if (--it->second == 0) {
_external_spill_directory_leases.erase(it);
}
}
-SpillDataDir* SpillFileManager::_get_store_for_spill() {
+SpillDataDir* SpillFileManager::_get_local_store_for_external_spill() {
+ // External writers use native filesystem paths, not Doris FileSystem
objects.
+ if (_remote_store != nullptr) {
+ return nullptr;
Review Comment:
[P2] Keep a local temporary root available to Paimon in S3 spill mode.
`ExecEnv::_init` creates only a `RemoteSpillDataDir` when
`spill_storage_type=s3`, and this return makes every external spill session
fail with `NO_AVAILABLE_ROOT_PATH`. `PaimonJniWriter` installs
`DorisIOManager`, whose first actual temporary-file request calls that session
and propagates an IO exception, so Paimon writes requiring temp files fail even
if local paths are configured. The new external-session test asserts this
failure; please retain a local root for Paimon or supply another supported
native path, with a Paimon temp-write test.
##########
fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java:
##########
@@ -563,6 +563,17 @@ public Set<Long> getExternalDmlAuditBackendIds() {
? Collections.emptySet() :
externalDmlAuditCoordinator.getDispatchedBackendIdsForAudit();
}
+ public Set<Long> getAuditStatisticsBackendIds() {
+ if (masterOpExecutor != null) {
+ return masterOpExecutor.getAuditStatisticsBackendIds();
+ }
+ Set<Long> backendIds =
Sets.newHashSet(getExternalDmlAuditBackendIds());
+ if (coord != null) {
Review Comment:
[P2] Preserve dispatched BE IDs when the coordinator is cleared after an
error. `executeAndSendResult` can dispatch fragments, then catch a
`getNext()`/result-send exception, cancel, and call `setCoord(null)` before
`handleQueryException` audits. This new method then returns no ordinary-query
participants, so `AuditLogHelper` submits without the final-report barrier; a
BE report delayed past `query_audit_log_timeout_ms` leaves zero or incomplete
remote spill bytes in the audit record. Retain the dispatched IDs independently
of mutable coordinator/result state (also relevant to forwarded cancellation),
and test an error after dispatch with a delayed final report. The earlier audit
thread covers the successful-query route.
##########
fe/fe-core/src/main/java/org/apache/doris/plugin/audit/AuditLoader.java:
##########
@@ -178,6 +178,8 @@ private void fillLogBuffer(AuditEvent event, StringBuilder
logBuffer) {
logBuffer.append(event.shuffleSendBytes).append(AUDIT_TABLE_COL_SEPARATOR);
logBuffer.append(event.spillWriteBytesToLocalStorage).append(AUDIT_TABLE_COL_SEPARATOR);
logBuffer.append(event.spillReadBytesFromLocalStorage).append(AUDIT_TABLE_COL_SEPARATOR);
+
logBuffer.append(event.spillWriteBytesToRemoteStorage).append(AUDIT_TABLE_COL_SEPARATOR);
Review Comment:
[P1] Keep audit rows loadable before the master expands the table schema.
During a follower-first rolling upgrade, a new follower emits these two values
and `AuditStreamLoader` names the new columns in every stream-load request,
while the old master has not run the master-only `InternalSchemaInitializer`.
The target table lacks those columns, so the load fails; `loadIfNecessary`
ignores the response and resets the batch, permanently losing audit rows from
the upgraded follower until the master changes the schema. Gate the new
row/column shape on the table schema or expand the schema first, retain failed
batches for retry, and cover this mixed-version sequence.
##########
be/src/exec/spill/spill_file_writer.cpp:
##########
@@ -90,42 +165,142 @@ 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 (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);
+ }
}
- if (_write_file_current_size) {
- COUNTER_UPDATE(_write_file_current_size, meta_size);
+
+ // Close synchronously. Spilling already runs on an IO thread that blocks
on every
+ // append, and with the default part size the wait at a part boundary (the
tail of the
+ // uploads plus one CompleteMultipartUpload) is negligible against the
part itself, so
+ // overlapping it with the next part is not worth an async close pipeline.
+ std::unique_ptr<io::FileWriter> writer = std::move(_file_writer);
+ MultipartUploadId upload = _multipart_upload_id(writer.get());
+ if (status.ok()) {
+ status = writer->close();
}
- _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() && writer->state() != io::FileWriter::State::CLOSED) {
+ // The part never reached a final state (the footer append failed).
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.
+ writer.reset();
}
- RETURN_IF_ERROR(_file_writer->close());
- _file_writer.reset();
+ _reconcile_part();
- // Advance to next part
- ++_current_part_index;
- if (spill_file) {
- spill_file->increment_part_count();
+ if (!status.ok()) {
+ LOG(WARNING) << "failed to close spill part " << _current_part_path <<
": " << status;
+ _abort_multipart_upload(upload);
+ } else if (spill_file) {
+ spill_file->add_part(_part_written_bytes);
Review Comment:
[P2] Count an object published before upload verification fails. With the
default-enabled S3 check, `PutObject` or `CompleteMultipartUpload` can succeed
and the following HEAD can fail. `close()` then returns an error and this
branch skips `add_part`, the only persisted-byte increment. If prefix cleanup
also fails, retry retains an actual object with `persisted_bytes=0`, so
`remote_spill_data_bytes()` and SHOW DATA omit it until deletion. Track
acknowledged or potentially published bytes separately from readable parts and
retain them through failed cleanup; test successful PUT/complete, failed HEAD,
and failed delete.
##########
be/src/exec/spill/spill_file.cpp:
##########
@@ -34,38 +34,64 @@
#include "util/debug_points.h"
namespace doris {
-SpillFile::SpillFile(SpillDataDir* data_dir, std::string relative_path)
+SpillFile::SpillFile(SpillDataDir* data_dir, io::FileSystemSPtr fs,
std::string relative_path)
: _data_dir(data_dir),
+ _fs(std::move(fs)),
_spill_dir(data_dir->get_spill_data_path() + "/" +
std::move(relative_path)) {}
+SpillFile::SpillFile(SpillDataDir* data_dir, std::string relative_path)
+ : SpillFile(data_dir, data_dir->fs(), std::move(relative_path)) {
+ DORIS_CHECK(!data_dir->is_remote());
+}
+
SpillFile::~SpillFile() {
gc();
}
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);
- 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());
- }
+ // The writer may outlive this file (e.g. when an error unwinds the file's
owner first).
+ // Drop its unfinished part now: closing it later would charge a footer to
no owner and
+ // could publish an object after the prefix below was deleted.
+ if (_active_writer != nullptr) {
+ _active_writer->_discard(this);
}
- // 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);
- _total_written_bytes = 0;
+ const int64_t written_bytes = std::exchange(_total_written_bytes, 0);
+ const int64_t persisted_bytes = std::exchange(_persisted_bytes, 0);
+ if (!_dir_created) {
+ _data_dir->release(written_bytes);
+ return;
+ }
+ _dir_created = false;
+ // 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.
+ // The filesystem is pinned when the file is created, not resolved from a
rotating default.
+ auto fs = _fs;
+ DORIS_CHECK(fs != nullptr) << "spill store " << _data_dir->path() << " is
not ready";
+ Status status = fs->delete_directory(_spill_dir);
+ DBUG_EXECUTE_IF("fault_inject::spill_file::gc", {
+ status = Status::Error<INTERNAL_ERROR>("fault_inject spill_file gc
failed");
+ });
+ if (status.ok()) {
+ _data_dir->release(written_bytes);
+ _data_dir->release_persisted_bytes(persisted_bytes);
+ return;
+ }
+ LOG_EVERY_T(WARNING, 1) << fmt::format("failed to delete spill data, dir
{}, error: {}",
+ _spill_dir, status.to_string());
+ // The data is still stored: keep it charged until a retry of the manager
deletes it.
+ auto* manager = ExecEnv::GetInstance()->spill_file_mgr();
+ if (manager == nullptr) {
+ _data_dir->release(written_bytes);
+ return;
+ }
+ manager->retry_spill_directory_deletion(_data_dir, _fs, _spill_dir,
written_bytes,
Review Comment:
[P2] Reconcile the reported stored bytes after a partial S3 delete.
`DeleteObjects` may delete some keys and report a per-key error, and recursive
deletion may fail after earlier batches succeeded. This failure path transfers
the full `persisted_bytes` to the retry queue; the manager subtracts it only
after an entirely successful retry. A persistently failing key therefore keeps
already deleted parts in SHOW DATA indefinitely. Retaining the full quota
reservation is conservative, but the stored-byte tally should reflect the
objects still present; test a partially successful multi-part delete.
--
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]