Coverage Report

Created: 2026-08-27 16:23

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/root/doris/be/src/io/io_common.h
Line
Count
Source
1
// Licensed to the Apache Software Foundation (ASF) under one
2
// or more contributor license agreements.  See the NOTICE file
3
// distributed with this work for additional information
4
// regarding copyright ownership.  The ASF licenses this file
5
// to you under the Apache License, Version 2.0 (the
6
// "License"); you may not use this file except in compliance
7
// with the License.  You may obtain a copy of the License at
8
//
9
//   http://www.apache.org/licenses/LICENSE-2.0
10
//
11
// Unless required by applicable law or agreed to in writing,
12
// software distributed under the License is distributed on an
13
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14
// KIND, either express or implied.  See the License for the
15
// specific language governing permissions and limitations
16
// under the License.
17
18
#pragma once
19
20
#include <gen_cpp/Types_types.h>
21
22
#include <optional>
23
#include <set>
24
#include <string>
25
26
namespace doris {
27
28
enum class ReaderType : uint8_t {
29
    READER_QUERY = 0,
30
    READER_ALTER_TABLE = 1,
31
    READER_BASE_COMPACTION = 2,
32
    READER_CUMULATIVE_COMPACTION = 3,
33
    READER_CHECKSUM = 4,
34
    READER_COLD_DATA_COMPACTION = 5,
35
    READER_SEGMENT_COMPACTION = 6,
36
    READER_FULL_COMPACTION = 7,
37
    UNKNOWN = 8
38
};
39
40
namespace io {
41
42
class RemoteScanCacheWriteLimiter;
43
enum class CacheWriteMode : uint8_t;
44
45
enum class FileCacheMissPolicy : uint8_t {
46
    READ_THROUGH_AND_WRITE_BACK = 0,
47
    REMOTE_ONLY_ON_MISS = 1,
48
};
49
50
struct FileReaderStats {
51
    size_t read_calls = 0;
52
    size_t read_bytes = 0;
53
    int64_t read_time_ns = 0;
54
    size_t read_rows = 0;
55
};
56
57
struct FileCacheStatistics {
58
    int64_t num_local_io_total = 0;
59
    int64_t num_remote_io_total = 0;
60
    int64_t num_peer_io_total = 0;
61
    int64_t local_io_timer = 0;
62
    int64_t bytes_read_from_local = 0;
63
    int64_t bytes_read_from_remote = 0;
64
    int64_t bytes_read_from_peer = 0;
65
    int64_t remote_io_timer = 0;
66
    int64_t peer_io_timer = 0;
67
    int64_t remote_wait_timer = 0;
68
    int64_t write_cache_io_timer = 0;
69
    int64_t bytes_write_into_cache = 0;
70
    int64_t num_skip_cache_io_total = 0;
71
    int64_t read_cache_file_directly_timer = 0;
72
    int64_t cache_get_or_set_timer = 0;
73
    int64_t lock_wait_timer = 0;
74
    int64_t get_timer = 0;
75
    int64_t set_timer = 0;
76
    int64_t async_cache_write_submitted = 0;
77
    int64_t async_cache_write_rejected = 0;
78
    int64_t async_cache_write_buffer_alloc_fail = 0;
79
    int64_t async_cache_write_drop_stale_epoch = 0;
80
    int64_t inflight_write_buffer_index_hit = 0;
81
    int64_t inflight_write_buffer_index_miss = 0;
82
    int64_t probe_downloaded_hit = 0;
83
    int64_t probe_downloading_hit = 0;
84
    int64_t probe_miss = 0;
85
    int64_t block_wait_success = 0;
86
    int64_t block_wait_timeout = 0;
87
88
    int64_t inverted_index_num_local_io_total = 0;
89
    int64_t inverted_index_num_remote_io_total = 0;
90
    int64_t inverted_index_num_peer_io_total = 0;
91
    int64_t inverted_index_bytes_read_from_local = 0;
92
    int64_t inverted_index_bytes_read_from_remote = 0;
93
    int64_t inverted_index_bytes_read_from_peer = 0;
94
    int64_t inverted_index_local_io_timer = 0;
95
    int64_t inverted_index_remote_io_timer = 0;
96
    int64_t inverted_index_peer_io_timer = 0;
97
    int64_t inverted_index_io_timer = 0;
98
    int64_t inverted_index_write_cache_io_timer = 0;
99
    int64_t inverted_index_bytes_write_into_cache = 0;
100
101
    int64_t segment_footer_index_num_local_io_total = 0;
102
    int64_t segment_footer_index_num_remote_io_total = 0;
103
    int64_t segment_footer_index_num_peer_io_total = 0;
104
    int64_t segment_footer_index_bytes_read_from_local = 0;
105
    int64_t segment_footer_index_bytes_read_from_remote = 0;
106
    int64_t segment_footer_index_bytes_read_from_peer = 0;
107
    int64_t segment_footer_index_local_io_timer = 0;
108
    int64_t segment_footer_index_remote_io_timer = 0;
109
    int64_t segment_footer_index_peer_io_timer = 0;
110
    int64_t segment_footer_index_write_cache_io_timer = 0;
111
    int64_t segment_footer_index_bytes_write_into_cache = 0;
112
    int64_t remote_only_on_miss_triggered = 0;
113
    int64_t remote_only_on_miss_threshold_bytes = 0;
114
115
    // Cross-CG / Same-CG peer read statistics
116
    int64_t num_cross_cg_peer_io_total = 0;
117
    int64_t bytes_read_from_cross_cg_peer = 0;
118
    int64_t cross_cg_peer_io_timer = 0; // nanoseconds
119
    int64_t num_same_cg_peer_io_total = 0;
120
    int64_t bytes_read_from_same_cg_peer = 0;
121
    int64_t same_cg_peer_io_timer = 0; // nanoseconds
122
    int64_t num_peer_race_peer_win = 0;
123
    int64_t num_peer_race_s3_win = 0;
124
    int64_t num_peer_lazy_fetch = 0;
125
    int64_t peer_lazy_fetch_timer = 0; // nanoseconds
126
127
    std::set<std::string> peer_hosts;
128
129
29
    void merge_from(const FileCacheStatistics& other) {
130
29
        num_local_io_total += other.num_local_io_total;
131
29
        num_remote_io_total += other.num_remote_io_total;
132
29
        num_peer_io_total += other.num_peer_io_total;
133
29
        local_io_timer += other.local_io_timer;
134
29
        bytes_read_from_local += other.bytes_read_from_local;
135
29
        bytes_read_from_remote += other.bytes_read_from_remote;
136
29
        bytes_read_from_peer += other.bytes_read_from_peer;
137
29
        remote_io_timer += other.remote_io_timer;
138
29
        peer_io_timer += other.peer_io_timer;
139
29
        remote_wait_timer += other.remote_wait_timer;
140
29
        write_cache_io_timer += other.write_cache_io_timer;
141
29
        bytes_write_into_cache += other.bytes_write_into_cache;
142
29
        num_skip_cache_io_total += other.num_skip_cache_io_total;
143
29
        read_cache_file_directly_timer += other.read_cache_file_directly_timer;
144
29
        cache_get_or_set_timer += other.cache_get_or_set_timer;
145
29
        lock_wait_timer += other.lock_wait_timer;
146
29
        get_timer += other.get_timer;
147
29
        set_timer += other.set_timer;
148
29
        async_cache_write_submitted += other.async_cache_write_submitted;
149
29
        async_cache_write_rejected += other.async_cache_write_rejected;
150
29
        async_cache_write_buffer_alloc_fail += other.async_cache_write_buffer_alloc_fail;
151
29
        async_cache_write_drop_stale_epoch += other.async_cache_write_drop_stale_epoch;
152
29
        inflight_write_buffer_index_hit += other.inflight_write_buffer_index_hit;
153
29
        inflight_write_buffer_index_miss += other.inflight_write_buffer_index_miss;
154
29
        probe_downloaded_hit += other.probe_downloaded_hit;
155
29
        probe_downloading_hit += other.probe_downloading_hit;
156
29
        probe_miss += other.probe_miss;
157
29
        block_wait_success += other.block_wait_success;
158
29
        block_wait_timeout += other.block_wait_timeout;
159
160
29
        inverted_index_num_local_io_total += other.inverted_index_num_local_io_total;
161
29
        inverted_index_num_remote_io_total += other.inverted_index_num_remote_io_total;
162
29
        inverted_index_num_peer_io_total += other.inverted_index_num_peer_io_total;
163
29
        inverted_index_bytes_read_from_local += other.inverted_index_bytes_read_from_local;
164
29
        inverted_index_bytes_read_from_remote += other.inverted_index_bytes_read_from_remote;
165
29
        inverted_index_bytes_read_from_peer += other.inverted_index_bytes_read_from_peer;
166
29
        inverted_index_local_io_timer += other.inverted_index_local_io_timer;
167
29
        inverted_index_remote_io_timer += other.inverted_index_remote_io_timer;
168
29
        inverted_index_peer_io_timer += other.inverted_index_peer_io_timer;
169
29
        inverted_index_io_timer += other.inverted_index_io_timer;
170
29
        inverted_index_write_cache_io_timer += other.inverted_index_write_cache_io_timer;
171
29
        inverted_index_bytes_write_into_cache += other.inverted_index_bytes_write_into_cache;
172
173
29
        segment_footer_index_num_local_io_total += other.segment_footer_index_num_local_io_total;
174
29
        segment_footer_index_num_remote_io_total += other.segment_footer_index_num_remote_io_total;
175
29
        segment_footer_index_num_peer_io_total += other.segment_footer_index_num_peer_io_total;
176
29
        segment_footer_index_bytes_read_from_local +=
177
29
                other.segment_footer_index_bytes_read_from_local;
178
29
        segment_footer_index_bytes_read_from_remote +=
179
29
                other.segment_footer_index_bytes_read_from_remote;
180
29
        segment_footer_index_bytes_read_from_peer +=
181
29
                other.segment_footer_index_bytes_read_from_peer;
182
29
        segment_footer_index_local_io_timer += other.segment_footer_index_local_io_timer;
183
29
        segment_footer_index_remote_io_timer += other.segment_footer_index_remote_io_timer;
184
29
        segment_footer_index_peer_io_timer += other.segment_footer_index_peer_io_timer;
185
29
        segment_footer_index_write_cache_io_timer +=
186
29
                other.segment_footer_index_write_cache_io_timer;
187
29
        segment_footer_index_bytes_write_into_cache +=
188
29
                other.segment_footer_index_bytes_write_into_cache;
189
29
        remote_only_on_miss_triggered =
190
29
                remote_only_on_miss_triggered || other.remote_only_on_miss_triggered;
191
29
        if (other.remote_only_on_miss_threshold_bytes > remote_only_on_miss_threshold_bytes) {
192
1
            remote_only_on_miss_threshold_bytes = other.remote_only_on_miss_threshold_bytes;
193
1
        }
194
195
29
        num_cross_cg_peer_io_total += other.num_cross_cg_peer_io_total;
196
29
        bytes_read_from_cross_cg_peer += other.bytes_read_from_cross_cg_peer;
197
29
        cross_cg_peer_io_timer += other.cross_cg_peer_io_timer;
198
29
        num_same_cg_peer_io_total += other.num_same_cg_peer_io_total;
199
29
        bytes_read_from_same_cg_peer += other.bytes_read_from_same_cg_peer;
200
29
        same_cg_peer_io_timer += other.same_cg_peer_io_timer;
201
29
        num_peer_race_peer_win += other.num_peer_race_peer_win;
202
29
        num_peer_race_s3_win += other.num_peer_race_s3_win;
203
29
        num_peer_lazy_fetch += other.num_peer_lazy_fetch;
204
29
        peer_lazy_fetch_timer += other.peer_lazy_fetch_timer;
205
206
29
        peer_hosts.insert(other.peer_hosts.begin(), other.peer_hosts.end());
207
29
    }
208
};
209
210
struct IOContext {
211
    ReaderType reader_type = ReaderType::UNKNOWN;
212
    // FIXME(plat1ko): Seems `is_disposable` can be inferred from the `reader_type`?
213
    bool is_disposable = false;
214
    bool is_index_data = false;
215
    bool read_file_cache = true;
216
    // TODO(lightman): use following member variables to control file cache
217
    bool is_persistent = false;
218
    // stop reader when reading, used in some interrupted operations
219
    bool should_stop = false;
220
    int64_t expiration_time = 0;
221
    const TUniqueId* query_id = nullptr;             // Ref
222
    FileCacheStatistics* file_cache_stats = nullptr; // Ref
223
    FileReaderStats* file_reader_stats = nullptr;    // Ref
224
    bool is_inverted_index = false;
225
    // if is_dryrun, read IO will download data to cache but return no data to reader
226
    // useful to skip cache data read from local disk to accelarate warm up
227
    bool is_dryrun = false;
228
    // if `is_warmup` == true, this I/O request is from a warm up task
229
    bool is_warmup {false};
230
    int64_t condition_cache_filtered_rows = 0;
231
    // Rows removed by file-local predicate conjuncts inside FileReader/TableReader. Scanner-level
232
    // output filtering already records its own unselected rows; this counter carries the rows that
233
    // were filtered before the block returned to Scanner.
234
    int64_t predicate_filtered_rows = 0;
235
    // if true, bypass peer read / peer-vs-S3 race and read directly from remote storage
236
    bool bypass_peer_read {false};
237
    // Per-call override for cache write completion semantics. An unset value follows the reader
238
    // option and the global async file-cache write switch.
239
    std::optional<CacheWriteMode> cache_write_mode_override = std::nullopt;
240
    FileCacheMissPolicy file_cache_miss_policy = FileCacheMissPolicy::READ_THROUGH_AND_WRITE_BACK;
241
    RemoteScanCacheWriteLimiter* remote_scan_cache_write_limiter = nullptr; // Ref
242
};
243
244
} // namespace io
245
} // namespace doris