Coverage Report

Created: 2026-10-07 08:42

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/io/fs/file_writer.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 <atomic>
21
#include <functional>
22
#include <future>
23
#include <memory>
24
25
#include "common/status.h"
26
#include "io/cache/block_file_cache.h"
27
#include "io/cache/block_file_cache_factory.h"
28
#include "io/cache/file_cache_common.h"
29
#include "io/fs/file_reader_writer_fwd.h"
30
#include "io/fs/path.h"
31
#include "util/slice.h"
32
33
namespace doris::io {
34
class FileSystem;
35
struct FileCacheAllocatorBuilder;
36
struct EncryptionInfo;
37
38
// Request statistics reported by remote file writers when the caller passes an instance
39
// through FileWriterOptions::remote_write_stats. All fields are cumulative. They count logical
40
// requests, one per call into the object storage client: attempts retried inside the client
41
// (throttling, transient errors) are not seen by the writer and not counted, so these are
42
// observability numbers, not the provider's billed request count.
43
struct RemoteWriteStats {
44
    std::atomic<int64_t> put_object_requests {0};
45
    std::atomic<int64_t> create_multipart_requests {0};
46
    std::atomic<int64_t> upload_part_requests {0};
47
    std::atomic<int64_t> complete_multipart_requests {0};
48
    std::atomic<int64_t> head_requests {0};
49
    std::atomic<int64_t> failed_requests {0};
50
    // Bytes acknowledged by object storage (PutObject and UploadPart payloads).
51
    std::atomic<int64_t> uploaded_bytes {0};
52
    // Sum of request latencies. Requests run concurrently, so this is not wall-clock time.
53
    std::atomic<int64_t> request_time_ns {0};
54
55
44
    int64_t total_requests() const {
56
44
        return put_object_requests + create_multipart_requests + upload_part_requests +
57
44
               complete_multipart_requests + head_requests;
58
44
    }
59
};
60
61
// Only affects remote file writers
62
struct FileWriterOptions {
63
    // S3 committer will start multipart uploading all files on BE side,
64
    // and then complete multipart upload these files on FE side.
65
    // If you do not complete multi parts of a file, the file will not be visible.
66
    // So in this way, the atomicity of a single file can be guaranteed. But it still cannot
67
    // guarantee the atomicity of multiple files.
68
    // Because hive committers have best-effort semantics,
69
    // this shortens the inconsistent time window.
70
    bool used_by_s3_committer = false;
71
    bool write_file_cache = false;
72
    bool allow_adaptive_file_cache_write = true;
73
    bool is_cold_data = false;
74
    bool sync_file_data = true;              // Whether flush data into storage system
75
    uint64_t file_cache_expiration_time = 0; // Absolute time, 0 means no TTL
76
    uint64_t approximate_bytes_to_write = 0; // Approximate bytes to write, used for file cache
77
    // Upload flow control, honoured by S3FileWriter only (other writers ignore both hooks).
78
    //
79
    // upload_submit_gate is called on the appending thread (appendv, or close for the last
80
    // buffer) right before a data buffer is submitted for upload, with the allocated capacity
81
    // of the buffer (s3_write_buffer_size, also for a partially filled last buffer) so that a
82
    // budget built on it bounds memory, not payload. It may block. A non-OK status fails the
83
    // writer: the buffer is dropped, no further data is accepted and close() reports the error.
84
    //
85
    // upload_done_callback is called exactly once, with the same capacity, for every buffer
86
    // that passed the gate — and only for those — when
87
    // the upload of that buffer has finished (success, provider error, or skipped because an
88
    // earlier buffer failed) and also when its submission failed. It runs on the upload thread
89
    // strictly before the buffer's status is published, so it always happens before the writer
90
    // reports a final close status or is destroyed. It must not block and must not touch the
91
    // FileWriter. Buffers that fail before being submitted (e.g. a checksum mismatch detected
92
    // inside the upload buffer) do not call back; callers that need exact accounting reconcile
93
    // after the writer reached its final state.
94
    std::function<Status(size_t)> upload_submit_gate = nullptr;
95
    std::function<void(size_t)> upload_done_callback = nullptr;
96
    // Optional sink for per-request statistics of remote file writers.
97
    std::shared_ptr<RemoteWriteStats> remote_write_stats = nullptr;
98
};
99
100
struct AsyncCloseStatusPack {
101
    std::promise<Status> promise;
102
    std::future<Status> future;
103
};
104
105
class FileWriter {
106
public:
107
    enum class State : uint8_t {
108
        OPENED = 0,
109
        ASYNC_CLOSING,
110
        CLOSED,
111
    };
112
273k
    FileWriter() = default;
113
268k
    virtual ~FileWriter() = default;
114
115
    FileWriter(const FileWriter&) = delete;
116
    const FileWriter& operator=(const FileWriter&) = delete;
117
118
    // Normal close. Wait for all data to persist before returning.
119
    // If there is no data appended, an empty file will be persisted.
120
    virtual Status close(bool non_block = false) = 0;
121
122
    // Non-blocking probe for a previous close(true).
123
    // OK means close finished successfully. NeedSendAgain means close is still running.
124
    // Other errors mean close finished with error or the writer does not support this API.
125
    // NOTE: This method consumes the async close result when it is ready. The caller must
126
    // use it as the only completion path for that async close; mixing it with close(false)
127
    // or another try_finish_close consumer is not supported.
128
0
    virtual Status try_finish_close() {
129
0
        return Status::NotSupported("try_finish_close is not supported");
130
0
    }
131
132
131k
    Status append(const Slice& data) { return appendv(&data, 1); }
133
134
    virtual Status appendv(const Slice* data, size_t data_cnt) = 0;
135
136
    virtual const Path& path() const = 0;
137
138
    virtual size_t bytes_appended() const = 0;
139
140
    virtual State state() const = 0;
141
142
    // Returns true if this file's data was written to a packed file.
143
    // Used to determine whether to collect packed slice location from PackedFileManager.
144
3
    virtual bool is_in_packed_file() const { return false; }
145
146
    FileCacheAllocatorBuilder* cache_builder() const {
147
        return _cache_builder == nullptr ? nullptr : _cache_builder.get();
148
    }
149
150
protected:
151
86.0k
    void init_cache_builder(const FileWriterOptions* opts, const Path& path) {
152
86.0k
        if (!config::enable_file_cache || opts == nullptr) {
153
1.70k
            return;
154
1.70k
        }
155
156
84.3k
        io::UInt128Wrapper path_hash = BlockFileCache::hash(path.filename().native());
157
84.3k
        BlockFileCache* file_cache_ptr = FileCacheFactory::instance()->get_by_path(path_hash);
158
159
84.3k
        bool has_enough_file_cache_space = opts->allow_adaptive_file_cache_write &&
160
84.3k
                                           config::enable_file_cache_adaptive_write &&
161
84.3k
                                           (opts->approximate_bytes_to_write > 0) &&
162
84.3k
                                           (file_cache_ptr->approximate_available_cache_size() >
163
6.24k
                                            opts->approximate_bytes_to_write);
164
165
84.3k
        VLOG_DEBUG << "path:" << path.filename().native()
166
66
                   << ", write_file_cache:" << opts->write_file_cache
167
66
                   << ", allow_adaptive_file_cache_write:" << opts->allow_adaptive_file_cache_write
168
66
                   << ", has_enough_file_cache_space:" << has_enough_file_cache_space
169
66
                   << ", approximate_bytes_to_write:" << opts->approximate_bytes_to_write
170
66
                   << ", file_cache_available_size:"
171
66
                   << file_cache_ptr->approximate_available_cache_size();
172
84.3k
        if (opts->write_file_cache || has_enough_file_cache_space) {
173
62.7k
            _cache_builder = std::make_unique<FileCacheAllocatorBuilder>(FileCacheAllocatorBuilder {
174
18.4E
                    opts ? opts->is_cold_data : false, opts ? opts->file_cache_expiration_time : 0,
175
62.7k
                    path_hash, file_cache_ptr});
176
62.7k
        }
177
84.3k
        return;
178
86.0k
    }
179
180
    std::unique_ptr<FileCacheAllocatorBuilder> _cache_builder =
181
            nullptr; // nullptr if disable write file cache
182
};
183
184
} // namespace doris::io