Coverage Report

Created: 2026-10-10 04:43

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exec/spill/spill_file_manager.h
Line
Count
Source
1
2
// Licensed to the Apache Software Foundation (ASF) under one
3
// or more contributor license agreements.  See the NOTICE file
4
// distributed with this work for additional information
5
// regarding copyright ownership.  The ASF licenses this file
6
// to you under the Apache License, Version 2.0 (the
7
// "License"); you may not use this file except in compliance
8
// with the License.  You may obtain a copy of the License at
9
//
10
//   http://www.apache.org/licenses/LICENSE-2.0
11
//
12
// Unless required by applicable law or agreed to in writing,
13
// software distributed under the License is distributed on an
14
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15
// KIND, either express or implied.  See the License for the
16
// specific language governing permissions and limitations
17
// under the License.
18
19
#pragma once
20
#include <atomic>
21
#include <memory>
22
#include <mutex>
23
#include <string>
24
#include <unordered_map>
25
#include <vector>
26
27
#include "common/metrics/metrics.h"
28
#include "common/status.h"
29
#include "exec/spill/spill_file.h"
30
#include "storage/options.h"
31
#include "util/stopwatch.hpp"
32
#include "util/threadpool.h"
33
34
namespace doris {
35
class RuntimeProfile;
36
template <typename T>
37
class AtomicCounter;
38
using IntAtomicCounter = AtomicCounter<int64_t>;
39
template <typename T>
40
class AtomicGauge;
41
using UIntGauge = AtomicGauge<uint64_t>;
42
class MetricEntity;
43
struct MetricPrototype;
44
class QueryContext;
45
class ResourceContext;
46
47
class SpillFileManager;
48
class SpillDataDir {
49
public:
50
    SpillDataDir(std::string path, int64_t capacity_bytes,
51
                 TStorageMedium::type storage_medium = TStorageMedium::HDD);
52
53
    Status init();
54
55
223
    const std::string& path() const { return _path; }
56
57
    std::string get_spill_data_path(const std::string& query_id = "") const;
58
59
    std::string get_spill_data_gc_path(const std::string& sub_dir_name = "") const;
60
61
662
    TStorageMedium::type storage_medium() const { return _storage_medium; }
62
63
    // check if the capacity reach the limit after adding the incoming data
64
    // return true if limit reached, otherwise, return false.
65
    bool reach_capacity_limit(int64_t incoming_data_size);
66
67
    Status update_capacity();
68
69
1.21k
    void update_spill_data_usage(int64_t incoming_data_size) {
70
1.21k
        std::lock_guard<std::mutex> l(_mutex);
71
1.21k
        _spill_data_bytes += incoming_data_size;
72
1.21k
        spill_disk_data_size->set_value(_spill_data_bytes);
73
1.21k
    }
74
75
11
    int64_t get_spill_data_bytes() {
76
11
        std::lock_guard<std::mutex> l(_mutex);
77
11
        return _spill_data_bytes;
78
11
    }
79
80
5
    int64_t get_spill_data_limit() {
81
5
        std::lock_guard<std::mutex> l(_mutex);
82
5
        return _spill_data_limit_bytes;
83
5
    }
84
85
    std::string debug_string();
86
87
private:
88
    bool _reach_disk_capacity_limit(int64_t incoming_data_size);
89
    void _update_inode_usage();
90
1.17k
    double _get_disk_usage(int64_t incoming_data_size) const {
91
1.17k
        return _disk_capacity_bytes == 0
92
1.17k
                       ? 0
93
1.17k
                       : (double)(_disk_capacity_bytes - _available_bytes + incoming_data_size) /
94
1.17k
                                 (double)_disk_capacity_bytes;
95
1.17k
    }
96
97
    friend class SpillFileManager;
98
    std::string _path;
99
100
    // protect _disk_capacity_bytes, _available_bytes, _spill_data_limit_bytes, _spill_data_bytes
101
    std::mutex _mutex;
102
    // the actual capacity of the disk of this data dir
103
    size_t _disk_capacity_bytes;
104
    int64_t _spill_data_limit_bytes = 0;
105
    // the actual available capacity of the disk of this data dir
106
    size_t _available_bytes = 0;
107
    int64_t _spill_data_bytes = 0;
108
    TStorageMedium::type _storage_medium;
109
110
    std::shared_ptr<MetricEntity> spill_data_dir_metric_entity;
111
    IntGauge* spill_disk_capacity = nullptr;
112
    IntGauge* spill_disk_limit = nullptr;
113
    IntGauge* spill_disk_avail_capacity = nullptr;
114
    IntGauge* spill_disk_data_size = nullptr;
115
    // inode statistics of the disk of this data dir, refreshed by update_capacity(). Spill creates
116
    // one directory per query and per operator plus one file per part, so inodes may be exhausted
117
    // long before bytes are. The gauges are the only storage of these values; the total is 0 when
118
    // the file system does not report a fixed inode count.
119
    IntGauge* spill_disk_inode_total = nullptr;
120
    IntGauge* spill_disk_inode_available = nullptr;
121
    // for test
122
    IntGauge* spill_disk_has_spill_data = nullptr;
123
    IntGauge* spill_disk_has_spill_gc_data = nullptr;
124
};
125
126
// Adapts one external writer to the same root selection, capacity accounting and query cleanup
127
// used by Doris spill files.
128
class ExternalSpillSession {
129
public:
130
    ~ExternalSpillSession();
131
132
    Status get_paths(std::vector<std::string>* paths);
133
134
    Status reserve(const std::string& path, int64_t bytes);
135
136
    void update_accounting(const std::string& path, int64_t current_bytes_delta,
137
                           int64_t write_bytes, int64_t read_bytes);
138
139
private:
140
    friend class SpillFileManager;
141
142
    ExternalSpillSession(SpillFileManager* manager, QueryContext* query_context,
143
                         std::string relative_path);
144
    bool _contains(const std::string& path) const;
145
146
    SpillFileManager* _manager;
147
    std::weak_ptr<QueryContext> _query_context;
148
    std::shared_ptr<ResourceContext> _resource_context;
149
    std::string _query_id;
150
    std::string _relative_path;
151
    SpillDataDir* _data_dir = nullptr;
152
    std::string _path;
153
    int64_t _accounted_bytes = 0;
154
    std::mutex _mutex;
155
};
156
157
class SpillFileManager {
158
public:
159
    ~SpillFileManager();
160
    SpillFileManager(
161
            std::unordered_map<std::string, std::unique_ptr<SpillDataDir>>&& spill_store_map);
162
163
    Status init();
164
165
    void stop();
166
167
    // Create SpillFile and register it
168
    // @param relative_path  Operator-formatted path under the spill root,
169
    //                       e.g. "query_id/sort-node_id-task_id-unique_id"
170
    Status create_spill_file(const std::string& relative_path, SpillFileSPtr& spill_file);
171
172
    // Create a lazy managed session for an external spill implementation. A spill root is selected
173
    // and registered only when the external implementation first requests its path.
174
    Status create_external_spill_session(const std::string& relative_path,
175
                                         QueryContext* query_context,
176
                                         std::unique_ptr<ExternalSpillSession>* spill_session);
177
178
    /// Get a unique ID for constructing spill file paths.
179
293
    uint64_t next_id() { return id_++; }
180
181
    // Delete SpillFile data synchronously.
182
    void delete_spill_file(SpillFileSPtr spill_file);
183
184
    // Recursively delete a per-query spill directory during query teardown. Failed deletions are
185
    // retained by the manager and retried by its GC and shutdown paths.
186
    void delete_query_spill_directory(const std::string& query_id, SpillDataDir* data_dir);
187
188
    void gc(int32_t max_work_time_ms);
189
190
723
    void update_spill_write_bytes(int64_t bytes) { _spill_write_bytes_counter->increment(bytes); }
191
192
330
    void update_spill_read_bytes(int64_t bytes) { _spill_read_bytes_counter->increment(bytes); }
193
194
private:
195
    friend class ExternalSpillSession;
196
197
    struct PendingQuerySpillDirectory {
198
        int failed_count {0};
199
        std::string query_dir;
200
    };
201
202
    // Per-store statistics of one gc round, logged in the gc summary.
203
    struct SpillGcStats {
204
        bool has_work = false;
205
        // Query directories still under the gc root after this round, i.e. the gc backlog.
206
        size_t backlog_dirs = 0;
207
        size_t deleted_dirs = 0;
208
        size_t deleted_files = 0;
209
        size_t failed_deletes = 0;
210
    };
211
212
    void _init_metrics();
213
    Status _init_spill_store_map();
214
    void _spill_gc_thread_callback();
215
    Status _try_delete_query_spill_directory(const PendingQuerySpillDirectory& pending_directory);
216
    void _retry_pending_query_spill_directories();
217
    Status _initialize_external_spill_session(ExternalSpillSession* spill_session);
218
    void _release_external_spill_session(ExternalSpillSession* spill_session);
219
220
    // Delete the gc backlog of one spill store until `max_work_time_ns` of `watch` has elapsed.
221
    void _gc_spill_store(SpillDataDir* store_dir, const MonotonicStopWatch& watch,
222
                         int64_t max_work_time_ns, SpillGcStats* stats);
223
    std::vector<SpillDataDir*> _get_stores_for_spill(TStorageMedium::type storage_medium);
224
    SpillDataDir* _get_store_for_spill();
225
226
    std::unordered_map<std::string, std::unique_ptr<SpillDataDir>> _spill_store_map;
227
228
    CountDownLatch _stop_background_threads_latch;
229
    std::shared_ptr<Thread> _spill_gc_thread;
230
231
    // Query cleanup uses the regular pending-deletion path. External leases only defer deletion
232
    // while an SDK task can still access the same query directory; filesystem I/O never holds this
233
    // mutex.
234
    std::mutex _pending_query_spill_directories_mutex;
235
    std::vector<PendingQuerySpillDirectory> _pending_query_spill_directories;
236
    std::unordered_map<std::string, size_t> _external_spill_directory_leases;
237
238
    std::atomic_uint64_t id_ = 0;
239
240
    std::shared_ptr<MetricEntity> _entity {nullptr};
241
242
    std::unique_ptr<doris::MetricPrototype> _spill_write_bytes_metric {nullptr};
243
    std::unique_ptr<doris::MetricPrototype> _spill_read_bytes_metric {nullptr};
244
245
    IntAtomicCounter* _spill_write_bytes_counter {nullptr};
246
    IntAtomicCounter* _spill_read_bytes_counter {nullptr};
247
};
248
} // namespace doris