github-actions[bot] commented on code in PR #68032:
URL: https://github.com/apache/doris/pull/68032#discussion_r4119558849


##########
be/src/exec/spill/spill_file_reader.cpp:
##########
@@ -62,77 +68,177 @@ SpillFileReader::SpillFileReader(RuntimeState* state, 
RuntimeProfile* profile,
     _read_file_size = get_counter(profile::SPILL_READ_FILE_BYTES);
     _read_rows_count = get_counter(profile::SPILL_READ_ROWS);
     _read_file_count = get_counter(profile::SPILL_READ_FILE_COUNT);
+    // Optional: older profiles may not register it.
+    _remote_read_requests = 
custom_profile->get_counter(profile::SPILL_REMOTE_READ_REQUESTS);
+    if (_is_remote) {
+        _coalesce_bytes = static_cast<size_t>(std::max<int64_t>(
+                0,
+                std::min(config::spill_s3_read_coalesce_bytes, 
state->spill_buffer_size_bytes())));
+    }
+}
+
+void SpillFileReader::_record_read(size_t bytes_read) {
+    COUNTER_UPDATE(_read_file_size, bytes_read);
+    
ExecEnv::GetInstance()->spill_file_mgr()->update_spill_read_bytes(bytes_read);
+    if (_is_remote) {
+        // One successful read_at() is one logical GET. Retries inside the 
object storage client
+        // and reads that finally failed are not counted (see 
RemoteWriteStats).
+        if (_remote_read_requests != nullptr) {
+            COUNTER_UPDATE(_remote_read_requests, 1);
+        }
+        if (_resource_ctx) {
+            
_resource_ctx->io_context()->update_spill_read_bytes_from_remote_storage(bytes_read);
+            _resource_ctx->io_context()->update_spill_remote_read_requests(1);
+        }
+        
ExecEnv::GetInstance()->spill_file_mgr()->update_spill_remote_read(bytes_read, 
1);
+    } else if (_resource_ctx) {
+        
_resource_ctx->io_context()->update_spill_read_bytes_from_local_storage(bytes_read);
+    }
 }
 
 Status SpillFileReader::open() {
     if (_is_open || _part_count == 0) {
         return Status::OK();
     }
-    RETURN_IF_ERROR(_open_part(0));
+    RETURN_IF_ERROR(_open_part(0, true));
     _is_open = true;
     return Status::OK();
 }
 
-Status SpillFileReader::_open_part(size_t part_index) {
+Status SpillFileReader::_open_part(size_t part_index, bool fetch_small_part) {
     _close_current_part();
 
     _current_part_index = part_index;
     _part_opened = true;
     std::string part_path = _spill_dir + "/" + std::to_string(part_index);
 
-    SCOPED_TIMER(_read_file_timer);
     COUNTER_UPDATE(_read_file_count, 1);
-    RETURN_IF_ERROR(io::global_local_filesystem()->open_file(part_path, 
&_file_reader));
-
-    size_t file_size = _file_reader->size();
-    DCHECK(file_size >= 16); // max_sub_block_size + block count
-
-    Slice result((char*)&_part_block_count, sizeof(size_t));
+    auto fs = _data_dir != nullptr ? _data_dir->fs() : 
io::global_local_filesystem();
+    if (fs == nullptr) {
+        return Status::InternalError("spill store {} is not ready", 
_data_dir->path());
+    }
+    io::FileReaderOptions opts;
+    opts.cache_type = io::FileCachePolicy::NO_CACHE;
+    // The writer recorded the part size; on object storage this saves a HEAD 
request.
+    opts.file_size = _part_sizes[part_index];
+    {
+        SCOPED_TIMER(_read_file_timer);
+        RETURN_IF_ERROR(fs->open_file(part_path, &_file_reader, &opts));
+    }
+    RETURN_IF_ERROR(_read_footer(_file_reader->size(), fetch_small_part));
+    _part_read_block_index = 0;
+    return Status::OK();
+}
 
-    // read block count
-    size_t bytes_read = 0;
-    RETURN_IF_ERROR(_file_reader->read_at(file_size - sizeof(size_t), result, 
&bytes_read));
-    DCHECK(bytes_read == 8);
-
-    // read max sub block size
-    bytes_read = 0;
-    result.data = (char*)&_part_max_sub_block_size;
-    RETURN_IF_ERROR(_file_reader->read_at(file_size - sizeof(size_t) * 2, 
result, &bytes_read));
-    DCHECK(bytes_read == 8);
-
-    // The buffer is used for two purposes:
-    // 1. Reading the block start offsets array (needs _part_block_count * 
sizeof(size_t) bytes)
-    // 2. Reading a single block's serialized data (needs up to 
_part_max_sub_block_size bytes)
-    // We must ensure the buffer is large enough for either case, so take the 
maximum.
-    size_t buff_size = std::max(_part_block_count * sizeof(size_t), 
_part_max_sub_block_size);
-    if (buff_size > _read_buff.size()) {
-        _read_buff.reserve(buff_size);
-    }
-
-    // Read the block start offsets array from the end of the file.
-    // The file layout (from end backwards) is:
+Status SpillFileReader::_read_footer(size_t file_size, bool fetch_small_part) {
+    // The part layout (from the end backwards) is:
     //   [block count (size_t)]
     //   [max sub block size (size_t)]
     //   [block start offsets array (_part_block_count * size_t)]
-    // So the offsets array starts at:
-    //   file_size - (_part_block_count + 2) * sizeof(size_t)
-    size_t read_offset = file_size - (_part_block_count + 2) * sizeof(size_t);
-    result.data = _read_buff.data();
-    result.size = _part_block_count * sizeof(size_t);
+    //   [serialized blocks]
+    constexpr size_t kFooterTailBytes = 2 * sizeof(size_t);
+    // Enough for the offsets of a few thousand blocks, so the footer usually 
takes one GET.
+    constexpr size_t kRemoteFooterProbeBytes = 64 * 1024;
+    if (file_size < kFooterTailBytes) {
+        return Status::InternalError("spill part {} is too small: {} bytes",
+                                     _file_reader->path().native(), file_size);
+    }
 
-    RETURN_IF_ERROR(_file_reader->read_at(read_offset, result, &bytes_read));
-    DCHECK(bytes_read == _part_block_count * sizeof(size_t));
+    // Without coalescing (local disk) the footer is read exactly. Otherwise a 
larger tail is
+    // read in one request, or the whole part when it fits in one coalesced 
read.
+    size_t probe_size = kFooterTailBytes;
+    if (_coalesce_bytes > 0) {
+        probe_size = fetch_small_part && file_size <= _coalesce_bytes

Review Comment:
   [P2] Apply the effective read cap to this footer probe. _coalesce_bytes is 
limited by spill_s3_read_coalesce_bytes and the query spill buffer, but a 1 KiB 
setting still reads a whole 40 KiB part here (or 64 KiB from a larger part). 
That violates the documented per-GET cap and allocates more buffer than the 
query requested. Bound the probe by that cap, allowing the 16-byte tail 
minimum, then fetch any remaining offsets separately; cover a small-cap part in 
a test.



##########
be/src/exec/spill/spill_file_manager.cpp:
##########
@@ -308,145 +456,8 @@ void SpillFileManager::gc(int32_t max_work_time_ms) {
     }
 }
 
-DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(spill_disk_capacity, MetricUnit::BYTES);
-DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(spill_disk_limit, MetricUnit::BYTES);
-DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(spill_disk_avail_capacity, 
MetricUnit::BYTES);
-DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(spill_disk_data_size, MetricUnit::BYTES);
-DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(spill_disk_has_spill_data, 
MetricUnit::BYTES);
-DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(spill_disk_has_spill_gc_data, 
MetricUnit::BYTES);
-
-SpillDataDir::SpillDataDir(std::string path, int64_t capacity_bytes,
-                           TStorageMedium::type storage_medium)
-        : _path(std::move(path)),
-          _disk_capacity_bytes(capacity_bytes),
-          _storage_medium(storage_medium) {
-    spill_data_dir_metric_entity = 
DorisMetrics::instance()->metric_registry()->register_entity(
-            std::string("spill_data_dir.") + _path, {{"path", _path + "/" + 
SPILL_DIR_PREFIX}});
-    INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity, 
spill_disk_capacity);
-    INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity, spill_disk_limit);
-    INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity, 
spill_disk_avail_capacity);
-    INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity, 
spill_disk_data_size);
-    INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity, 
spill_disk_has_spill_data);
-    INT_GAUGE_METRIC_REGISTER(spill_data_dir_metric_entity, 
spill_disk_has_spill_gc_data);
-}
-
-bool is_directory_empty(const std::filesystem::path& dir) {
-    // Spill cleanup may delete the directory while the iterator is 
constructed or advanced. Treat
-    // that race as empty for these presence metrics.
-    try {
-        return std::filesystem::is_directory(dir) &&
-               std::filesystem::directory_iterator(dir) ==
-                       
std::filesystem::end(std::filesystem::directory_iterator {});
-    } catch (const std::filesystem::filesystem_error&) {
-        return true;
-    }
-}
-
-Status SpillDataDir::init() {
-    bool exists = false;
-    RETURN_IF_ERROR(io::global_local_filesystem()->exists(_path, &exists));
-    if (!exists) {
-        RETURN_NOT_OK_STATUS_WITH_WARN(Status::IOError("opendir failed, 
path={}", _path),
-                                       "check file exist failed");
-    }
-    RETURN_IF_ERROR(update_capacity());
-    LOG(INFO) << fmt::format(
-            "spill storage path: {}, capacity: {}, limit: {}, available: "
-            "{}",
-            _path, PrettyPrinter::print_bytes(_disk_capacity_bytes),
-            PrettyPrinter::print_bytes(_spill_data_limit_bytes),
-            PrettyPrinter::print_bytes(_available_bytes));
-    return Status::OK();
-}
-
-std::string SpillDataDir::get_spill_data_path(const std::string& query_id) 
const {
-    auto dir = fmt::format("{}/{}", _path, SPILL_DIR_PREFIX);
-    if (!query_id.empty()) {
-        dir = fmt::format("{}/{}", dir, query_id);
-    }
-    return dir;
-}
-
-std::string SpillDataDir::get_spill_data_gc_path(const std::string& 
sub_dir_name) const {
-    auto dir = fmt::format("{}/{}", _path, SPILL_GC_DIR_PREFIX);
-    if (!sub_dir_name.empty()) {
-        dir = fmt::format("{}/{}", dir, sub_dir_name);
-    }
-    return dir;
+int64_t SpillFileManager::remote_spill_data_bytes() {

Review Comment:
   [P2] Separate stored bytes from admission reservations for this billing 
value. _write_internal reserves each block before S3FileWriter submits its 
pending buffer, so an open writer can keep a nonzero _spill_data_bytes entirely 
in memory. This method sends that value to SHOW DATA as bytes held in object 
storage, producing a fresh billable size with no object. Keep the conservative 
reservation for the capacity limit, but report acknowledged object-storage 
bytes separately; test an open writer before close.



##########
fe/fe-core/src/main/java/org/apache/doris/common/proc/CurrentQueryStatisticsProcDir.java:
##########
@@ -45,7 +45,9 @@ public class CurrentQueryStatisticsProcDir implements 
ProcDirInterface {
             .add("ScanBytesFromLocalStorage").add("ScanBytesFromRemoteStorage")
             
.add("SpillWriteBytesToLocalStorage").add("SpillReadBytesFromLocalStorage")
             .add("BytesWriteIntoCache")
-            .add("TotalTasks").add("FinishedTasks").add("Progress").build();
+            .add("TotalTasks").add("FinishedTasks").add("Progress")
+            // Appended last: multi-FE aggregation concatenates rows by 
position.

Review Comment:
   [P2] Normalize remote FE rows during a rolling upgrade. These two titles add 
two cells, but QueryProfileAction.currentQueries concatenates every FE's 
NodeInfo.rows and labels them with the serving FE's TITLE_NAMES without using 
each response's columnNames. A new FE therefore returns old-FE rows two cells 
short of its advertised schema (and an old FE can return unlabeled extra 
cells). Map or pad rows by each response's headers and test mixed FE versions.



##########
fe/fe-core/src/main/java/org/apache/doris/qe/AuditLogHelper.java:
##########
@@ -274,6 +274,10 @@ private static void logAuditLogImpl(ConnectContext ctx, 
String origStmt, Stateme
                         statistics.getSpillWriteBytesToLocalStorage())
                 .setSpillReadBytesFromLocalStorage(statistics == null ? 0 :
                         statistics.getSpillReadBytesFromLocalStorage())
+                // Remote spill bytes are only reported through 
TQueryStatistics; for queries

Review Comment:
   [P2] Seed remote spill from final query statistics too. This hardcodes zero 
while PQueryStatistics supplies the local spill pair but has no remote fields. 
Ordinary audit events are released after query_audit_log_timeout_ms without 
waiting for a final BE TQueryStatistics report, so a delayed report or longer 
BE reporting interval persists a completed S3-spilling query with zero remote 
bytes. Carry the pair in final PQueryStatistics or wait for final 
participating-BE reports, and test a report delayed past the audit timeout.



##########
be/src/exec/spill/remote_spill_data_dir.cpp:
##########
@@ -0,0 +1,133 @@
+// 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/remote_spill_data_dir.h"
+
+#include <glog/logging.h>
+
+#include <utility>
+
+#include "cloud/cloud_storage_engine.h"
+#include "cloud/config.h"
+#include "common/config.h"
+#include "common/logging.h"
+#include "common/metrics/metrics.h"
+#include "io/fs/remote_file_system.h"
+#include "runtime/exec_env.h"
+#include "service/backend_options.h"
+#include "storage/olap_define.h"
+#include "storage/storage_policy.h"
+#include "util/pretty_printer.h"
+
+namespace doris {
+
+RemoteSpillDataDir::RemoteSpillDataDir(std::string vault_id)
+        : SpillDataDir(fmt::format("s3:{}", vault_id.empty() ? "default" : 
vault_id),
+                       /*spill_root=*/"",
+                       fmt::format("s3:{}", vault_id.empty() ? "default" : 
vault_id),
+                       /*capacity_bytes=*/0, TStorageMedium::S3),
+          _vault_id(std::move(vault_id)) {}
+
+Status RemoteSpillDataDir::init() {
+    RETURN_IF_ERROR(update_capacity());
+    LOG(INFO) << fmt::format("remote spill store registered, vault_id={}, 
limit={}",
+                             _vault_id.empty() ? "<default>" : _vault_id,
+                             
PrettyPrinter::print_bytes(_spill_data_limit_bytes));
+    return Status::OK();
+}
+
+Status RemoteSpillDataDir::ensure_ready() {
+    if (ready()) {

Review Comment:
   [P1] Re-resolve the default vault for new spill files. With 
spill_s3_storage_vault empty, the first spill binds this store to default A; 
after SET DEFAULT STORAGE VAULT B and a BE refresh, this ready() return keeps 
every later file on A. If A is retired, new spill fails although B is 
available. Resolve the current default for each new file and pin its filesystem 
through read and GC, preserving existing A files; test an A-to-B rotation with 
both files live.



##########
fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/RemoteSpillStatsPoller.java:
##########
@@ -0,0 +1,161 @@
+// 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.
+
+package org.apache.doris.cloud.catalog;
+
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.AnalysisException;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.Pair;
+import org.apache.doris.common.Status;
+import org.apache.doris.common.util.MasterDaemon;
+import org.apache.doris.proto.InternalService;
+import org.apache.doris.rpc.BackendServiceProxy;
+import org.apache.doris.system.Backend;
+import org.apache.doris.thrift.TStatusCode;
+
+import com.google.common.annotations.VisibleForTesting;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * Polls the bytes of query spill held in object storage 
(spill_storage_type=s3) from the alive
+ * backends of all clusters, so that SHOW DATA reads them from memory. Runs on 
every FE on its own
+ * schedule (cloud_spill_stats_poll_interval_second): the freshness of the 
value must not depend on
+ * how long a round of the tablet stats takes.
+ */
+public class RemoteSpillStatsPoller extends MasterDaemon {
+    private static final Logger LOG = 
LogManager.getLogger(RemoteSpillStatsPoller.class);
+
+    private static final int RPC_TIMEOUT_SECOND = 5;
+
+    /**
+     * One successful poll: the value and when it was fetched, on the 
monotonic clock so that a wall
+     * clock moved backward cannot keep a stale value fresh.
+     */
+    private static final class RemoteSpillStats {
+        private final long bytes;
+        private final long fetchTimeNanos;
+
+        private RemoteSpillStats(long bytes, long fetchTimeNanos) {
+            this.bytes = bytes;
+            this.fetchTimeNanos = fetchTimeNanos;
+        }
+    }
+
+    // Summed over the alive BEs of all clusters. Null until the first 
successful poll. A BE that is
+    // gone no longer contributes.
+    private volatile RemoteSpillStats remoteSpillStats = null;
+
+    public RemoteSpillStatsPoller() {
+        super("remote spill stats poller", pollIntervalMs());
+    }
+
+    private static long pollIntervalMs() {
+        return Math.max(1, Config.cloud_spill_stats_poll_interval_second) * 
1000L;
+    }
+
+    /**
+     * A value older than this is not served: the configured max age, but at 
least three poll
+     * intervals so that a longer interval cannot make every value stale.
+     */
+    @VisibleForTesting
+    static long maxAgeSecond() {
+        return Math.max(Config.cloud_spill_stats_max_age_second,
+                3L * Math.max(1, 
Config.cloud_spill_stats_poll_interval_second));
+    }
+
+    @Override
+    protected void runAfterCatalogReady() {
+        refresh();
+        // The interval is mutable.
+        setInterval(pollIntervalMs());
+    }
+
+    private void refresh() {
+        List<Backend> backends;
+        try {
+            backends = 
Env.getCurrentSystemInfo().getAllBackendsByAllCluster().values().asList();
+        } catch (AnalysisException e) {
+            LOG.warn("failed to list the backends for the remote spill stats", 
e);
+            return;
+        }
+        InternalService.PGetBeResourceRequest request = 
InternalService.PGetBeResourceRequest.newBuilder().build();
+        List<Pair<Backend, Future<InternalService.PGetBeResourceResponse>>> 
futures = new ArrayList<>();
+        for (Backend be : backends) {
+            if (!be.isAlive()) {

Review Comment:
   [P2] Poll reachable registered BEs before publishing a fresh billing total. 
A failed heartbeat marks a BE not alive, but heartbeat uses a different port 
from this get_be_resource BRPC call. A running BE with a blocked heartbeat port 
can therefore still return its S3 spill bytes, yet this skip excludes them and 
refresh() timestamps the lower sum as current. Try the stats RPC for registered 
BEs with failed heartbeats and include successful replies; test heartbeat 
failure with a reachable BRPC endpoint.



-- 
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