mrhhsg commented on code in PR #68032:
URL: https://github.com/apache/doris/pull/68032#discussion_r4119385586
##########
be/src/exec/spill/spill_file_writer.cpp:
##########
@@ -90,42 +142,136 @@ 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();
Review Comment:
Fixed in bdd4fad: `_do_spill()` and `_do_intermediate_merge()` return
`writer->close()`'s status. The intermediate merge closes its output before
deleting its inputs.
--
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]