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 |