mrhhsg commented on code in PR #68032:
URL: https://github.com/apache/doris/pull/68032#discussion_r4229735722
##########
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:
已在 3bccf51 修复:S3 模式保留本地 spill root(未配置时回退到 storage_root_path),查询 spill 仍只选
S3,Paimon 外部 session
只选本地路径。ExternalSpillSessionUsesLocalRootWhileQuerySpillUsesS3 测试覆盖 native
临时文件创建/写入及并行的 S3 查询 spill。
##########
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:
已在 3bccf51 修复:失败 close
不再等同于未发布;等待上传结束后先保守计入候选字节,再列举该文件前缀核对实际对象,清理失败时将该值交给重试队列。HEAD 校验失败现在返回错误状态而非触发
DCHECK。PublishedPartRemainsBillableWhenVerificationAndCleanupFail 覆盖 PUT
成功、HEAD 失败和删除失败。
##########
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:
已在 3bccf51 修复:删除失败后对文件前缀重新 list,并在初次 GC
及后续重试中把已存储字节调到实际剩余对象大小;容量预留仍等到完全删除后释放。PartialDeleteReconcilesOnlyRemainingObjectBytes
使用首批删除成功、后续批次失败的 mock 验证。
##########
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:
已在 3bccf51 修复:executeAndSendResult 的 finally 在 coordinator 关闭前原子保存本次 query
ID 和已 dispatch 的 BE 集合,清空 coordinator 后审计仍可等待最终报告;转发 cancel 覆盖 result 时也保留原查询
BE ID。新增 post-dispatch 错误后延迟最终报告测试,以及转发取消和重试 query ID 测试。
##########
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:
已在 3bccf51 修复:audit stream load 根据当前 audit_log 表 schema 选择旧/新列形状,旧 master
尚未加列时从批次中投影掉两个新字段;失败批次保留并重试,传输结果不确定时复用唯一 label,避免重复写入,且不接受过滤行。测试覆盖旧/新及部分
schema、失败重试和 label 已完成响应。
--
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]