Coverage Report

Created: 2026-10-07 08:42

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exec/spill/spill_file_reader.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/data.pb.h>
21
22
#include <memory>
23
#include <string>
24
#include <vector>
25
26
#include "common/status.h"
27
#include "core/pod_array.h"
28
#include "core/pod_array_fwd.h"
29
#include "io/fs/file_reader_writer_fwd.h"
30
#include "runtime/runtime_profile.h"
31
#include "runtime/workload_management/resource_context.h"
32
#include "util/slice.h"
33
34
namespace doris {
35
class RuntimeState;
36
class Block;
37
class SpillDataDir;
38
39
/// SpillFileReader reads blocks sequentially across all parts of a SpillFile.
40
///
41
/// Usage:
42
///   auto reader = spill_file->create_reader(state, profile);
43
///   RETURN_IF_ERROR(reader->open());
44
///   bool eos = false;
45
///   while (!eos) { RETURN_IF_ERROR(reader->read(&block, &eos)); }
46
///
47
/// Part boundaries are transparent to the caller. When the current part is
48
/// exhausted, the reader automatically opens the next part.
49
///
50
/// Parts are opened on the SpillDataDir's file system (local disk or object storage).
51
/// Part sizes are known from the writer, so no size lookup is needed on open.
52
///
53
/// On object storage every read is a GET, so the reader keeps the request count low:
54
/// the part footer is fetched with one tail read of at most `_coalesce_bytes` (plus one more
55
/// read for the rest of a block offset array that does not fit), adjacent blocks are coalesced
56
/// into one read of at most `_coalesce_bytes`, and a part no larger than that is fetched whole.
57
/// A single block, or the rest of the offset array, larger than `_coalesce_bytes` is read whole.
58
class SpillFileReader {
59
public:
60
    SpillFileReader(RuntimeState* state, RuntimeProfile* profile, SpillDataDir* data_dir,
61
                    std::string spill_dir, std::vector<int64_t> part_sizes);
62
63
195
    ~SpillFileReader() { (void)close(); }
64
65
    /// Open the first part and read its footer metadata.
66
    Status open();
67
68
    /// Read the next block. Automatically advances across part boundaries.
69
    /// Sets *eos = true when all parts are exhausted.
70
    Status read(Block* block, bool* eos);
71
72
    /// Seek to a global block index within the whole spill file.
73
    /// block_index is 0-based across all parts.
74
    /// If block_index is out of range, the reader is positioned at EOS.
75
    Status seek(size_t block_index);
76
77
    Status close();
78
79
private:
80
    /// Open a specific part file and read its footer. With `fetch_small_part`, a part that
81
    /// fits in one coalesced read is fetched whole, so its blocks need no further reads.
82
    Status _open_part(size_t part_index, bool fetch_small_part);
83
84
    /// Read the footer (block offsets, max sub block size, block count) of the current part.
85
    Status _read_footer(size_t file_size, bool fetch_small_part);
86
87
    /// Point `out` at the serialized bytes of block `index` of the current part, reading
88
    /// them (and the following blocks that fit in the coalesce window) if not buffered.
89
    Status _block_slice(size_t index, Slice* out);
90
91
    /// Read exactly `len` bytes at `offset` of the current part into `data`.
92
    Status _read_exact(size_t offset, char* data, size_t len);
93
94
    /// Seek implementation with status propagation.
95
    Status _seek_to_block(size_t block_index);
96
97
    /// Close the current part's file reader.
98
    void _close_current_part();
99
100
    /// Make the read buffer hold at least `size` bytes.
101
    void _ensure_read_buff(size_t size);
102
103
    /// Account bytes read from the store (local or remote) as one request.
104
    void _record_read(size_t bytes_read);
105
106
    // ── Configuration ──
107
    SpillDataDir* _data_dir = nullptr;
108
    std::string _spill_dir;
109
    std::vector<int64_t> _part_sizes;
110
    size_t _part_count;
111
    bool _is_remote = false;
112
    // Upper bound of one coalesced read; 0 reads block by block (local disk).
113
    size_t _coalesce_bytes = 0;
114
115
    // ── Current part state ──
116
    size_t _current_part_index = 0;
117
    bool _is_open = false;
118
    bool _part_opened = false;
119
    io::FileReaderSPtr _file_reader;
120
    size_t _part_block_count = 0;
121
    size_t _part_read_block_index = 0;
122
    size_t _part_max_sub_block_size = 0;
123
    // Holds the bytes of [_window_begin, _window_end) of the current part.
124
    PaddedPODArray<char> _read_buff;
125
    size_t _window_begin = 0;
126
    size_t _window_end = 0;
127
    std::vector<size_t> _block_start_offsets;
128
129
    PBlock _pb_block;
130
131
    // ── Counters ──
132
    RuntimeProfile::Counter* _read_file_timer = nullptr;
133
    RuntimeProfile::Counter* _deserialize_timer = nullptr;
134
    RuntimeProfile::Counter* _read_block_count = nullptr;
135
    RuntimeProfile::Counter* _read_block_data_size = nullptr;
136
    RuntimeProfile::Counter* _read_file_size = nullptr;
137
    RuntimeProfile::Counter* _read_rows_count = nullptr;
138
    RuntimeProfile::Counter* _read_file_count = nullptr;
139
    // Remote only, may be null when the profile does not register it.
140
    RuntimeProfile::Counter* _remote_read_requests = nullptr;
141
142
    std::shared_ptr<ResourceContext> _resource_ctx = nullptr;
143
};
144
145
using SpillFileReaderSPtr = std::shared_ptr<SpillFileReader>;
146
147
} // namespace doris