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]

Reply via email to