/src/rocksdb/db/external_sst_file_ingestion_job.cc
Line | Count | Source |
1 | | // Copyright (c) 2011-present, Facebook, Inc. All rights reserved. |
2 | | // This source code is licensed under both the GPLv2 (found in the |
3 | | // COPYING file in the root directory) and Apache 2.0 License |
4 | | // (found in the LICENSE.Apache file in the root directory). |
5 | | |
6 | | #include "db/external_sst_file_ingestion_job.h" |
7 | | |
8 | | #include <algorithm> |
9 | | #include <cinttypes> |
10 | | #include <string> |
11 | | #include <unordered_set> |
12 | | #include <vector> |
13 | | |
14 | | #include "db/builder.h" |
15 | | #include "db/db_impl/db_impl.h" |
16 | | #include "db/version_edit.h" |
17 | | #include "file/file_util.h" |
18 | | #include "file/random_access_file_reader.h" |
19 | | #include "logging/logging.h" |
20 | | #include "monitoring/statistics_impl.h" |
21 | | #include "options/options_helper.h" |
22 | | #include "table/merging_iterator.h" |
23 | | #include "table/prepared_file_info.h" |
24 | | #include "table/sst_file_writer_collectors.h" |
25 | | #include "table/table_builder.h" |
26 | | #include "table/unique_id_impl.h" |
27 | | #include "test_util/sync_point.h" |
28 | | #include "util/udt_util.h" |
29 | | |
30 | | namespace ROCKSDB_NAMESPACE { |
31 | | |
32 | 0 | bool ExternalSstFileIngestionJob::SupportsAtomicReplaceRangeTombstone() const { |
33 | 0 | return cfd_->is_delete_range_supported() && |
34 | 0 | cfd_->ioptions().compaction_style == kCompactionStyleUniversal && |
35 | 0 | ingestion_options_.allow_global_seqno && |
36 | 0 | !ingestion_options_.allow_db_generated_files && |
37 | 0 | !ingestion_options_.snapshot_consistency && |
38 | 0 | ucmp_->timestamp_size() == 0; |
39 | 0 | } |
40 | | |
41 | | bool ExternalSstFileIngestionJob::CanUseAtomicReplaceRangeTombstone( |
42 | 0 | const std::optional<RangeOpt>& atomic_replace_range) const { |
43 | 0 | return atomic_replace_range.has_value() && |
44 | 0 | atomic_replace_range->start.has_value() && |
45 | 0 | atomic_replace_range->limit.has_value() && |
46 | 0 | SupportsAtomicReplaceRangeTombstone() && |
47 | 0 | !ingestion_options_.fail_if_not_bottommost_level; |
48 | 0 | } |
49 | | |
50 | | bool ExternalSstFileIngestionJob::HasPartialOverlap( |
51 | 0 | const VersionStorageInfo* vstorage) const { |
52 | 0 | assert(vstorage != nullptr); |
53 | 0 | assert(atomic_replace_range_.has_value()); |
54 | 0 | assert(!atomic_replace_range_->unset()); |
55 | 0 | for (int level = 0; level < cfd_->NumberLevels(); ++level) { |
56 | 0 | for (const auto* file : vstorage->LevelFiles(level)) { |
57 | 0 | if (file_range_checker_.Overlaps(*atomic_replace_range_, file->smallest, |
58 | 0 | file->largest) && |
59 | 0 | !file_range_checker_.Contains(*atomic_replace_range_, file->smallest, |
60 | 0 | file->largest)) { |
61 | 0 | return true; |
62 | 0 | } |
63 | 0 | } |
64 | 0 | } |
65 | 0 | return false; |
66 | 0 | } |
67 | | |
68 | | size_t ExternalSstFileIngestionJob::NumFilesToPrepare( |
69 | | size_t num_external_files, |
70 | 0 | const std::optional<RangeOpt>& atomic_replace_range) const { |
71 | 0 | return num_external_files + |
72 | 0 | static_cast<size_t>( |
73 | 0 | CanUseAtomicReplaceRangeTombstone(atomic_replace_range)); |
74 | 0 | } |
75 | | |
76 | | Status ExternalSstFileIngestionJob::PrepareAtomicReplaceRangeTombstone( |
77 | | const Slice& start, const Slice& limit, uint64_t file_number, |
78 | 0 | SuperVersion* super_version) { |
79 | 0 | assert(cfd_->is_delete_range_supported()); |
80 | 0 | assert(cfd_->ioptions().compaction_style == kCompactionStyleUniversal); |
81 | 0 | assert(ingestion_options_.allow_global_seqno); |
82 | 0 | assert(!ingestion_options_.allow_db_generated_files); |
83 | 0 | assert(!ingestion_options_.fail_if_not_bottommost_level); |
84 | 0 | assert(!ingestion_options_.snapshot_consistency); |
85 | 0 | assert(ucmp_->timestamp_size() == 0); |
86 | |
|
87 | 0 | const uint32_t path_id = 0; |
88 | 0 | const std::string file_path = |
89 | 0 | TableFileName(cfd_->ioptions().cf_paths, file_number, path_id); |
90 | 0 | ExternalSstFileInfo external_file_info; |
91 | 0 | Status status; |
92 | 0 | { |
93 | 0 | Options options( |
94 | 0 | BuildDBOptions(db_options_, mutable_db_options_), |
95 | 0 | BuildColumnFamilyOptions(cfd_->initial_cf_options(), |
96 | 0 | super_version->mutable_cf_options)); |
97 | 0 | SstFileWriter writer(env_options_, options, ucmp_); |
98 | 0 | status = writer.Open( |
99 | 0 | file_path, super_version->mutable_cf_options.default_write_temperature); |
100 | 0 | if (status.ok()) { |
101 | 0 | status = writer.DeleteRange(start, limit); |
102 | 0 | } |
103 | 0 | if (status.ok()) { |
104 | 0 | status = writer.Finish(&external_file_info); |
105 | 0 | } |
106 | 0 | } |
107 | 0 | if (!status.ok()) { |
108 | 0 | fs_->DeleteFile(file_path, IOOptions(), nullptr).PermitUncheckedError(); |
109 | 0 | return status; |
110 | 0 | } |
111 | | |
112 | 0 | IngestedFileInfo tombstone; |
113 | 0 | tombstone.file_temperature = |
114 | 0 | super_version->mutable_cf_options.default_write_temperature; |
115 | 0 | tombstone.prefetch_lmax_index_and_filter_blocks = |
116 | 0 | ingestion_options_.prefetch_lmax_index_and_filter_blocks; |
117 | 0 | const PreparedFileInfo* prepared_file_info = |
118 | 0 | ingestion_options_.write_global_seqno |
119 | 0 | ? nullptr |
120 | 0 | : external_file_info.prepared_file_info.get(); |
121 | 0 | status = GetIngestedFileInfo(file_path, file_number, prepared_file_info, |
122 | 0 | &tombstone, super_version); |
123 | 0 | if (!status.ok()) { |
124 | 0 | fs_->DeleteFile(file_path, IOOptions(), nullptr).PermitUncheckedError(); |
125 | 0 | return status; |
126 | 0 | } |
127 | | |
128 | 0 | tombstone.internal_file_path = file_path; |
129 | 0 | tombstone.copy_file = true; |
130 | 0 | tombstone.generated_for_ingestion = true; |
131 | 0 | tombstone.file_checksum = external_file_info.file_checksum; |
132 | 0 | tombstone.file_checksum_func_name = |
133 | 0 | external_file_info.file_checksum_func_name; |
134 | 0 | atomic_replace_range_tombstone_.emplace(std::move(tombstone)); |
135 | 0 | return Status::OK(); |
136 | 0 | } |
137 | | |
138 | 0 | void ExternalSstFileIngestionJob::ActivateAtomicReplaceRangeTombstone() { |
139 | 0 | assert(atomic_replace_range_tombstone_.has_value()); |
140 | | // Prepare rejects an empty external file list and verifies every input file |
141 | | // is contained in the replacement range. The covering tombstone must |
142 | | // therefore overlap at least one input file. |
143 | 0 | assert(!files_to_ingest_.empty()); |
144 | 0 | file_batches_to_ingest_.clear(); |
145 | 0 | files_to_ingest_.emplace_back(std::move(*atomic_replace_range_tombstone_)); |
146 | 0 | std::rotate(files_to_ingest_.begin(), files_to_ingest_.end() - 1, |
147 | 0 | files_to_ingest_.end()); |
148 | 0 | atomic_replace_range_tombstone_.reset(); |
149 | 0 | atomic_replace_range_tombstone_active_ = true; |
150 | 0 | files_overlap_ = ComputeFilesOverlap(files_to_ingest_); |
151 | 0 | assert(files_overlap_); |
152 | 0 | DivideInputFilesIntoBatches(); |
153 | 0 | } |
154 | | |
155 | | Status ExternalSstFileIngestionJob::Prepare( |
156 | | const std::vector<std::string>& external_files_paths, |
157 | | const std::vector<std::string>& files_checksums, |
158 | | const std::vector<std::string>& files_checksum_func_names, |
159 | | const std::vector<const PreparedFileInfo*>& file_infos, |
160 | | const std::optional<RangeOpt>& atomic_replace_range, |
161 | | const Temperature& file_temperature, uint64_t next_file_number, |
162 | 0 | SuperVersion* sv) { |
163 | 0 | Status status; |
164 | | |
165 | | // Read the information of files we are ingesting. |
166 | 0 | for (size_t i = 0; i < external_files_paths.size(); i++) { |
167 | 0 | const std::string& file_path = external_files_paths[i]; |
168 | 0 | IngestedFileInfo file_to_ingest; |
169 | | // For temperature, first assume it matches provided hint |
170 | 0 | file_to_ingest.file_temperature = file_temperature; |
171 | 0 | file_to_ingest.prefetch_lmax_index_and_filter_blocks = |
172 | 0 | ingestion_options_.prefetch_lmax_index_and_filter_blocks; |
173 | 0 | const PreparedFileInfo* prepared_file_info = |
174 | 0 | file_infos.empty() ? nullptr : file_infos[i]; |
175 | 0 | status = GetIngestedFileInfo(file_path, next_file_number++, |
176 | 0 | prepared_file_info, &file_to_ingest, sv); |
177 | 0 | if (!status.ok()) { |
178 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
179 | 0 | "Failed to get ingested file info: %s: %s", |
180 | 0 | file_path.c_str(), status.ToString().c_str()); |
181 | 0 | return status; |
182 | 0 | } |
183 | | |
184 | | // Files generated in another DB or CF may have a different column family |
185 | | // ID, so we let it pass here. |
186 | 0 | if (file_to_ingest.cf_id != |
187 | 0 | TablePropertiesCollectorFactory::Context::kUnknownColumnFamily && |
188 | 0 | file_to_ingest.cf_id != cfd_->GetID() && |
189 | 0 | !ingestion_options_.allow_db_generated_files) { |
190 | 0 | return Status::InvalidArgument( |
191 | 0 | "External file column family id don't match"); |
192 | 0 | } |
193 | | |
194 | 0 | if (file_to_ingest.num_entries == 0 && |
195 | 0 | file_to_ingest.num_range_deletions == 0) { |
196 | 0 | return Status::InvalidArgument("File contain no entries"); |
197 | 0 | } |
198 | | |
199 | 0 | if (!file_to_ingest.smallest_internal_key.Valid() || |
200 | 0 | !file_to_ingest.largest_internal_key.Valid()) { |
201 | 0 | return Status::Corruption("Generated table have corrupted keys"); |
202 | 0 | } |
203 | | |
204 | 0 | files_to_ingest_.emplace_back(std::move(file_to_ingest)); |
205 | 0 | } |
206 | | |
207 | 0 | auto num_files = files_to_ingest_.size(); |
208 | 0 | if (num_files == 0) { |
209 | 0 | return Status::InvalidArgument("The list of files is empty"); |
210 | 0 | } |
211 | | // Detect whether the input files overlap one another; this drives how they |
212 | | // are divided into batches below. |
213 | 0 | files_overlap_ = ComputeFilesOverlap(files_to_ingest_); |
214 | |
|
215 | 0 | if (atomic_replace_range.has_value()) { |
216 | 0 | atomic_replace_range_.emplace(); |
217 | |
|
218 | 0 | if (atomic_replace_range->start && atomic_replace_range->limit) { |
219 | 0 | if (ucmp_->CompareWithoutTimestamp(*atomic_replace_range->start, |
220 | 0 | *atomic_replace_range->limit) >= 0) { |
221 | 0 | return Status::InvalidArgument( |
222 | 0 | "Atomic replace range limit must be greater than start"); |
223 | 0 | } |
224 | | // User keys to internal keys (with timestamps) |
225 | 0 | const size_t ts_sz = ucmp_->timestamp_size(); |
226 | 0 | std::string start_with_ts, limit_with_ts; |
227 | 0 | auto [start, limit] = MaybeAddTimestampsToRange( |
228 | 0 | atomic_replace_range->start, atomic_replace_range->limit, ts_sz, |
229 | 0 | &start_with_ts, &limit_with_ts); |
230 | 0 | assert(start.has_value()); |
231 | 0 | assert(limit.has_value()); |
232 | 0 | atomic_replace_range_->smallest_internal_key.Set( |
233 | 0 | *start, kMaxSequenceNumber, kValueTypeForSeek); |
234 | 0 | atomic_replace_range_->largest_internal_key = |
235 | 0 | RangeTombstone(*start, *limit, kMaxSequenceNumber).SerializeEndKey(); |
236 | | // Check files to ingest against replace range |
237 | 0 | for (size_t i = 0; i < num_files; i++) { |
238 | 0 | if (!file_range_checker_.Contains(*atomic_replace_range_, |
239 | 0 | files_to_ingest_[i])) { |
240 | 0 | return Status::InvalidArgument( |
241 | 0 | "Atomic replace range does not contain all files"); |
242 | 0 | } |
243 | 0 | } |
244 | 0 | } else { |
245 | | // Currently if either bound is not present, both must be |
246 | 0 | assert(atomic_replace_range->start.has_value() == false); |
247 | 0 | assert(atomic_replace_range->limit.has_value() == false); |
248 | 0 | assert(atomic_replace_range_->smallest_internal_key.unset()); |
249 | 0 | assert(atomic_replace_range_->largest_internal_key.unset()); |
250 | 0 | } |
251 | 0 | } |
252 | | |
253 | 0 | if (files_overlap_) { |
254 | 0 | if (ingestion_options_.ingest_behind) { |
255 | 0 | return Status::NotSupported( |
256 | 0 | "Files with overlapping ranges cannot be ingested with ingestion " |
257 | 0 | "behind mode."); |
258 | 0 | } |
259 | | |
260 | | // Overlapping files need at least two different sequence numbers. If |
261 | | // settings disables global seqno, ingestion will fail anyway, so fail |
262 | | // fast in prepare. |
263 | 0 | if (!ingestion_options_.allow_global_seqno && |
264 | 0 | !ingestion_options_.allow_db_generated_files) { |
265 | 0 | return Status::InvalidArgument( |
266 | 0 | "Global seqno is required, but disabled (because external files key " |
267 | 0 | "range overlaps)."); |
268 | 0 | } |
269 | | |
270 | 0 | if (ucmp_->timestamp_size() > 0) { |
271 | 0 | return Status::NotSupported( |
272 | 0 | "Files with overlapping ranges cannot be ingested to column " |
273 | 0 | "family with user-defined timestamp enabled."); |
274 | 0 | } |
275 | 0 | } |
276 | | |
277 | | // Copy/Move external files into DB |
278 | 0 | std::unordered_set<size_t> ingestion_path_ids; |
279 | 0 | for (IngestedFileInfo& f : files_to_ingest_) { |
280 | 0 | f.copy_file = false; |
281 | 0 | const std::string path_outside_db = f.external_file_path; |
282 | 0 | const std::string path_inside_db = TableFileName( |
283 | 0 | cfd_->ioptions().cf_paths, f.fd.GetNumber(), f.fd.GetPathId()); |
284 | 0 | if (ingestion_options_.move_files || ingestion_options_.link_files) { |
285 | 0 | status = |
286 | 0 | fs_->LinkFile(path_outside_db, path_inside_db, IOOptions(), nullptr); |
287 | 0 | if (status.ok()) { |
288 | | // It is unsafe to assume application had sync the file and file |
289 | | // directory before ingest the file. For integrity of RocksDB we need |
290 | | // to sync the file. |
291 | 0 | TEST_SYNC_POINT("ExternalSstFileIngestionJob::BeforeSyncIngestedFile"); |
292 | 0 | Status s = fs_->SyncFile(path_inside_db, env_options_, IOOptions(), |
293 | 0 | db_options_.use_fsync, nullptr); |
294 | 0 | TEST_SYNC_POINT("ExternalSstFileIngestionJob::AfterSyncIngestedFile"); |
295 | 0 | TEST_SYNC_POINT_CALLBACK( |
296 | 0 | "ExternalSstFileIngestionJob::CheckSyncReturnCode", &s); |
297 | | // Some file systems (especially remote/distributed) don't support |
298 | | // explicitly syncing the file and don't require it. Ignore the |
299 | | // NotSupported error in that case. |
300 | 0 | if (!s.IsNotSupported()) { |
301 | 0 | status = s; |
302 | 0 | if (!status.ok()) { |
303 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
304 | 0 | "Failed to sync ingested file %s: %s", |
305 | 0 | path_inside_db.c_str(), status.ToString().c_str()); |
306 | 0 | } |
307 | 0 | } |
308 | 0 | } else if (status.IsNotSupported() && |
309 | 0 | ingestion_options_.failed_move_fall_back_to_copy) { |
310 | | // Original file is on a different FS, use copy instead of hard linking. |
311 | 0 | f.copy_file = true; |
312 | 0 | ROCKS_LOG_INFO(db_options_.info_log, |
313 | 0 | "Tried to link file %s but it's not supported : %s", |
314 | 0 | path_outside_db.c_str(), status.ToString().c_str()); |
315 | 0 | } else { |
316 | 0 | ROCKS_LOG_WARN(db_options_.info_log, "Failed to link file %s to %s: %s", |
317 | 0 | path_outside_db.c_str(), path_inside_db.c_str(), |
318 | 0 | status.ToString().c_str()); |
319 | 0 | } |
320 | 0 | } else { |
321 | 0 | f.copy_file = true; |
322 | 0 | } |
323 | |
|
324 | 0 | if (f.copy_file) { |
325 | 0 | TEST_SYNC_POINT_CALLBACK("ExternalSstFileIngestionJob::Prepare:CopyFile", |
326 | 0 | nullptr); |
327 | | // Always determining the destination temperature from the ingested-to |
328 | | // level would be difficult because in general we only find out the level |
329 | | // ingested to later, during Run(). |
330 | | // However, we can guarantee "last level" temperature for when the user |
331 | | // requires ingestion to the last level. |
332 | 0 | Temperature dst_temp = |
333 | 0 | (ingestion_options_.ingest_behind || |
334 | 0 | ingestion_options_.fail_if_not_bottommost_level) |
335 | 0 | ? sv->mutable_cf_options.last_level_temperature |
336 | 0 | : sv->mutable_cf_options.default_write_temperature; |
337 | | // Note: CopyFile also syncs the new file. |
338 | 0 | status = CopyFile(fs_.get(), path_outside_db, f.file_temperature, |
339 | 0 | path_inside_db, dst_temp, 0, db_options_.use_fsync, |
340 | 0 | io_tracer_); |
341 | | // The destination of the copy will be ingested |
342 | 0 | f.file_temperature = dst_temp; |
343 | |
|
344 | 0 | if (!status.ok()) { |
345 | 0 | ROCKS_LOG_WARN(db_options_.info_log, "Failed to copy file %s to %s: %s", |
346 | 0 | path_outside_db.c_str(), path_inside_db.c_str(), |
347 | 0 | status.ToString().c_str()); |
348 | 0 | } |
349 | 0 | } else { |
350 | | // Note: we currently assume that linking files does not cross |
351 | | // temperatures, so no need to change f.file_temperature |
352 | 0 | } |
353 | 0 | TEST_SYNC_POINT("ExternalSstFileIngestionJob::Prepare:FileAdded"); |
354 | 0 | if (!status.ok()) { |
355 | 0 | break; |
356 | 0 | } |
357 | 0 | f.internal_file_path = path_inside_db; |
358 | | // Initialize the checksum information of ingested files. |
359 | 0 | f.file_checksum = kUnknownFileChecksum; |
360 | 0 | f.file_checksum_func_name = kUnknownFileChecksumFuncName; |
361 | 0 | ingestion_path_ids.insert(f.fd.GetPathId()); |
362 | 0 | } |
363 | |
|
364 | 0 | if (status.ok() && CanUseAtomicReplaceRangeTombstone(atomic_replace_range) && |
365 | 0 | HasPartialOverlap(sv->current->storage_info())) { |
366 | 0 | assert(atomic_replace_range->start.has_value()); |
367 | 0 | assert(atomic_replace_range->limit.has_value()); |
368 | 0 | status = PrepareAtomicReplaceRangeTombstone(*atomic_replace_range->start, |
369 | 0 | *atomic_replace_range->limit, |
370 | 0 | next_file_number, sv); |
371 | 0 | if (status.ok()) { |
372 | 0 | assert(atomic_replace_range_tombstone_.has_value()); |
373 | 0 | ingestion_path_ids.insert( |
374 | 0 | atomic_replace_range_tombstone_->fd.GetPathId()); |
375 | 0 | } |
376 | 0 | } |
377 | |
|
378 | 0 | TEST_SYNC_POINT("ExternalSstFileIngestionJob::BeforeSyncDir"); |
379 | 0 | if (status.ok()) { |
380 | 0 | for (auto path_id : ingestion_path_ids) { |
381 | 0 | status = directories_->GetDataDir(path_id)->FsyncWithDirOptions( |
382 | 0 | IOOptions(), nullptr, |
383 | 0 | DirFsyncOptions(DirFsyncOptions::FsyncReason::kNewFileSynced)); |
384 | 0 | if (!status.ok()) { |
385 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
386 | 0 | "Failed to sync directory %" ROCKSDB_PRIszt |
387 | 0 | " while ingest file: %s", |
388 | 0 | path_id, status.ToString().c_str()); |
389 | 0 | break; |
390 | 0 | } |
391 | 0 | } |
392 | 0 | } |
393 | 0 | TEST_SYNC_POINT("ExternalSstFileIngestionJob::AfterSyncDir"); |
394 | | |
395 | | // Generate and check the sst file checksum. Note that, if |
396 | | // IngestExternalFileOptions::write_global_seqno is true, we will not update |
397 | | // the checksum information in the files_to_ingests_ here, since the file is |
398 | | // updated with the new global_seqno. After global_seqno is updated, DB will |
399 | | // generate the new checksum and store it in the Manifest. In all other cases |
400 | | // if ingestion_options_.write_global_seqno == true and |
401 | | // verify_file_checksum is false, we only check the checksum function name. |
402 | 0 | if (status.ok() && db_options_.file_checksum_gen_factory != nullptr) { |
403 | 0 | if (ingestion_options_.verify_file_checksum == false && |
404 | 0 | files_checksums.size() == files_to_ingest_.size() && |
405 | 0 | files_checksum_func_names.size() == files_to_ingest_.size()) { |
406 | | // Only when verify_file_checksum == false and the checksum for ingested |
407 | | // files are provided, DB will use the provided checksum and does not |
408 | | // generate the checksum for ingested files. |
409 | 0 | need_generate_file_checksum_ = false; |
410 | 0 | } else { |
411 | 0 | need_generate_file_checksum_ = true; |
412 | 0 | } |
413 | 0 | std::vector<std::string> generated_checksums; |
414 | 0 | std::vector<std::string> generated_checksum_func_names; |
415 | | // Step 1: generate the checksum for ingested sst file. |
416 | 0 | if (need_generate_file_checksum_) { |
417 | 0 | for (size_t i = 0; i < files_to_ingest_.size(); i++) { |
418 | 0 | std::string generated_checksum; |
419 | 0 | std::string generated_checksum_func_name; |
420 | 0 | std::string requested_checksum_func_name = |
421 | 0 | i < files_checksum_func_names.size() ? files_checksum_func_names[i] |
422 | 0 | : ""; |
423 | | // TODO: rate limit file reads for checksum calculation during file |
424 | | // ingestion. |
425 | | // TODO: plumb Env::IOActivity |
426 | 0 | ReadOptions ro; |
427 | | // Pass user-provided checksums through FileOptions when available. |
428 | | // The caller may not have provided checksums at all (empty vectors), |
429 | | // so we guard with a bounds check. |
430 | 0 | FileOptions fopts; |
431 | 0 | if (i < files_checksums.size()) { |
432 | 0 | fopts.file_checksum = files_checksums[i]; |
433 | 0 | } |
434 | 0 | if (i < files_checksum_func_names.size()) { |
435 | 0 | fopts.file_checksum_func_name = files_checksum_func_names[i]; |
436 | 0 | } else { |
437 | 0 | fopts.file_checksum_func_name = kNoFileChecksumFuncName; |
438 | 0 | } |
439 | 0 | IOStatus io_s = GenerateOneFileChecksum( |
440 | 0 | fs_.get(), files_to_ingest_[i].internal_file_path, |
441 | 0 | db_options_.file_checksum_gen_factory.get(), |
442 | 0 | requested_checksum_func_name, &generated_checksum, |
443 | 0 | &generated_checksum_func_name, |
444 | 0 | ingestion_options_.verify_checksums_readahead_size, |
445 | 0 | db_options_.allow_mmap_reads, io_tracer_, |
446 | 0 | db_options_.rate_limiter.get(), ro, db_options_.stats, |
447 | 0 | db_options_.clock, fopts); |
448 | 0 | if (!io_s.ok()) { |
449 | 0 | status = io_s; |
450 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
451 | 0 | "Sst file checksum generation of file: %s failed: %s", |
452 | 0 | files_to_ingest_[i].internal_file_path.c_str(), |
453 | 0 | status.ToString().c_str()); |
454 | 0 | break; |
455 | 0 | } |
456 | 0 | if (ingestion_options_.write_global_seqno == false) { |
457 | 0 | files_to_ingest_[i].file_checksum = generated_checksum; |
458 | 0 | files_to_ingest_[i].file_checksum_func_name = |
459 | 0 | generated_checksum_func_name; |
460 | 0 | } |
461 | 0 | generated_checksums.push_back(generated_checksum); |
462 | 0 | generated_checksum_func_names.push_back(generated_checksum_func_name); |
463 | 0 | } |
464 | 0 | } |
465 | | |
466 | | // Step 2: based on the verify_file_checksum and ingested checksum |
467 | | // information, do the verification. |
468 | 0 | if (status.ok()) { |
469 | 0 | if (files_checksums.size() == files_to_ingest_.size() && |
470 | 0 | files_checksum_func_names.size() == files_to_ingest_.size()) { |
471 | | // Verify the checksum and checksum function name. |
472 | 0 | if (ingestion_options_.verify_file_checksum) { |
473 | 0 | for (size_t i = 0; i < files_to_ingest_.size(); i++) { |
474 | 0 | if (files_checksum_func_names[i] != |
475 | 0 | generated_checksum_func_names[i]) { |
476 | 0 | status = Status::InvalidArgument( |
477 | 0 | "DB file checksum gen factory " + |
478 | 0 | std::string(db_options_.file_checksum_gen_factory->Name()) + |
479 | 0 | " generated checksum function name " + |
480 | 0 | generated_checksum_func_names[i] + " for file " + |
481 | 0 | external_files_paths[i] + |
482 | 0 | " which does not match requested/provided " + |
483 | 0 | files_checksum_func_names[i]); |
484 | 0 | break; |
485 | 0 | } |
486 | 0 | if (files_checksums[i] != generated_checksums[i]) { |
487 | 0 | status = Status::Corruption( |
488 | 0 | "Checksum verification mismatch for ingestion file " + |
489 | 0 | external_files_paths[i] + " using function " + |
490 | 0 | generated_checksum_func_names[i] + ". Expected: " + |
491 | 0 | Slice(files_checksums[i]).ToString(/*hex=*/true) + |
492 | 0 | " Computed: " + |
493 | 0 | Slice(generated_checksums[i]).ToString(/*hex=*/true)); |
494 | 0 | break; |
495 | 0 | } |
496 | 0 | } |
497 | 0 | } else { |
498 | | // If verify_file_checksum is not enabled, we only verify the factory |
499 | | // recognizes the checksum function name. If it does not match, fail |
500 | | // the ingestion. If matches, we trust the ingested checksum |
501 | | // information and store in the Manifest. |
502 | 0 | for (size_t i = 0; i < files_to_ingest_.size(); i++) { |
503 | 0 | FileChecksumGenContext gen_context; |
504 | 0 | gen_context.file_name = files_to_ingest_[i].internal_file_path; |
505 | 0 | gen_context.requested_checksum_func_name = |
506 | 0 | files_checksum_func_names[i]; |
507 | 0 | auto file_checksum_gen = |
508 | 0 | db_options_.file_checksum_gen_factory |
509 | 0 | ->CreateFileChecksumGenerator(gen_context); |
510 | |
|
511 | 0 | if (file_checksum_gen == nullptr || |
512 | 0 | files_checksum_func_names[i] != file_checksum_gen->Name()) { |
513 | 0 | status = Status::InvalidArgument( |
514 | 0 | "Checksum function name " + files_checksum_func_names[i] + |
515 | 0 | " for file " + external_files_paths[i] + |
516 | 0 | " not recognized by DB checksum gen factory" + |
517 | 0 | db_options_.file_checksum_gen_factory->Name() + |
518 | 0 | (file_checksum_gen ? (" Returned function " + |
519 | 0 | std::string(file_checksum_gen->Name())) |
520 | 0 | : "")); |
521 | 0 | break; |
522 | 0 | } |
523 | 0 | files_to_ingest_[i].file_checksum = files_checksums[i]; |
524 | 0 | files_to_ingest_[i].file_checksum_func_name = |
525 | 0 | files_checksum_func_names[i]; |
526 | 0 | } |
527 | 0 | } |
528 | 0 | } else if (files_checksums.size() != files_checksum_func_names.size() || |
529 | 0 | files_checksums.size() != 0) { |
530 | | // The checksum or checksum function name vector are not both empty |
531 | | // and they are incomplete. |
532 | 0 | status = Status::InvalidArgument( |
533 | 0 | "The checksum information of ingested sst files are nonempty and " |
534 | 0 | "the size of checksums or the size of the checksum function " |
535 | 0 | "names does not match with the number of ingested sst files"); |
536 | 0 | } |
537 | 0 | if (!status.ok()) { |
538 | 0 | ROCKS_LOG_WARN(db_options_.info_log, "Ingestion failed: %s", |
539 | 0 | status.ToString().c_str()); |
540 | 0 | } |
541 | 0 | } |
542 | 0 | } |
543 | |
|
544 | 0 | if (status.ok()) { |
545 | 0 | DivideInputFilesIntoBatches(); |
546 | 0 | } |
547 | |
|
548 | 0 | return status; |
549 | 0 | } |
550 | | |
551 | 0 | void ExternalSstFileIngestionJob::DivideInputFilesIntoBatches() { |
552 | 0 | if (!files_overlap_) { |
553 | | // No overlap, treat as one batch without the need of tracking overall batch |
554 | | // range. |
555 | 0 | file_batches_to_ingest_.emplace_back(/* _track_batch_range= */ false); |
556 | 0 | for (auto& file : files_to_ingest_) { |
557 | 0 | file_batches_to_ingest_.back().AddFile(&file, file_range_checker_); |
558 | 0 | } |
559 | 0 | return; |
560 | 0 | } |
561 | | |
562 | 0 | file_batches_to_ingest_.emplace_back(/* _track_batch_range= */ true); |
563 | 0 | for (auto& file : files_to_ingest_) { |
564 | 0 | if (!file_batches_to_ingest_.back().unset() && |
565 | 0 | file_range_checker_.Overlaps(file_batches_to_ingest_.back(), file, |
566 | 0 | /* known_sorted= */ false)) { |
567 | 0 | file_batches_to_ingest_.emplace_back(/* _track_batch_range= */ true); |
568 | 0 | } |
569 | 0 | file_batches_to_ingest_.back().AddFile(&file, file_range_checker_); |
570 | 0 | } |
571 | 0 | } |
572 | | |
573 | | bool ExternalSstFileIngestionJob::ComputeFilesOverlap( |
574 | 0 | const autovector<IngestedFileInfo>& files) const { |
575 | 0 | const size_t num_files = files.size(); |
576 | 0 | if (num_files <= 1) { |
577 | 0 | return false; |
578 | 0 | } |
579 | | // Verify whether the files have overlapping ranges by sorting copies of the |
580 | | // file ranges and checking adjacent pairs. |
581 | 0 | autovector<const IngestedFileInfo*> sorted_files; |
582 | 0 | for (size_t i = 0; i < num_files; i++) { |
583 | 0 | sorted_files.push_back(&files[i]); |
584 | 0 | } |
585 | 0 | std::sort(sorted_files.begin(), sorted_files.end(), file_range_checker_); |
586 | 0 | for (size_t i = 0; i + 1 < num_files; i++) { |
587 | 0 | if (file_range_checker_.Overlaps(*sorted_files[i], *sorted_files[i + 1], |
588 | 0 | /* known_sorted= */ true)) { |
589 | 0 | return true; |
590 | 0 | } |
591 | 0 | } |
592 | 0 | return false; |
593 | 0 | } |
594 | | |
595 | | Status ExternalSstFileIngestionJob::MergeForSameColumnFamily( |
596 | 0 | ExternalSstFileIngestionJob* other) { |
597 | 0 | assert(other != nullptr); |
598 | 0 | assert(other != this); |
599 | 0 | assert(cfd_ == other->cfd_); |
600 | 0 | if (atomic_replace_range_.has_value() || |
601 | 0 | other->atomic_replace_range_.has_value()) { |
602 | 0 | return Status::NotSupported( |
603 | 0 | "cannot merge file ingestion handles for the same column family when " |
604 | 0 | "atomic_replace_range is used"); |
605 | 0 | } |
606 | 0 | if (!(ingestion_options_ == other->ingestion_options_)) { |
607 | 0 | return Status::InvalidArgument( |
608 | 0 | "file ingestion handles for the same column family must be prepared " |
609 | 0 | "with the same IngestExternalFileOptions"); |
610 | 0 | } |
611 | | // Append the other job's prepared files after this job's so that, for any |
612 | | // overlapping keys, the other job's data wins via a higher assigned sequence |
613 | | // number -- the same semantics as passing all the files to a single ingestion |
614 | | // call in this order. Recompute overlap and rebuild the batches over the |
615 | | // union. |
616 | 0 | for (IngestedFileInfo& file : other->files_to_ingest_) { |
617 | 0 | files_to_ingest_.push_back(std::move(file)); |
618 | 0 | } |
619 | 0 | other->files_to_ingest_.clear(); |
620 | 0 | other->file_batches_to_ingest_.clear(); |
621 | 0 | files_overlap_ = ComputeFilesOverlap(files_to_ingest_); |
622 | 0 | file_batches_to_ingest_.clear(); |
623 | 0 | DivideInputFilesIntoBatches(); |
624 | 0 | return Status::OK(); |
625 | 0 | } |
626 | | |
627 | | Status ExternalSstFileIngestionJob::NeedsFlush(bool* flush_needed, |
628 | 0 | SuperVersion* super_version) { |
629 | 0 | Status status; |
630 | 0 | if (atomic_replace_range_.has_value() && atomic_replace_range_->unset()) { |
631 | | // For replacing whole CF, we can simply check whether memtable is empty |
632 | 0 | *flush_needed = !super_version->mem->IsEmpty(); |
633 | 0 | } else { |
634 | 0 | autovector<UserKeyRange> ranges; |
635 | 0 | if (atomic_replace_range_.has_value()) { |
636 | 0 | assert(!atomic_replace_range_->smallest_internal_key.unset()); |
637 | 0 | assert(!atomic_replace_range_->largest_internal_key.unset()); |
638 | | // NOTE: we already checked in Prepare() that the atomic_replace_range |
639 | | // covers all the files_to_ingest. |
640 | 0 | ranges.emplace_back( |
641 | 0 | atomic_replace_range_->smallest_internal_key.user_key(), |
642 | 0 | atomic_replace_range_->largest_internal_key.user_key()); |
643 | 0 | } else { |
644 | 0 | ranges.reserve(files_to_ingest_.size()); |
645 | 0 | for (const IngestedFileInfo& file_to_ingest : files_to_ingest_) { |
646 | 0 | ranges.emplace_back(file_to_ingest.start_ukey, |
647 | 0 | file_to_ingest.limit_ukey); |
648 | 0 | } |
649 | 0 | } |
650 | 0 | status = cfd_->RangesOverlapWithMemtables( |
651 | 0 | ranges, super_version, db_options_.allow_data_in_errors, flush_needed, |
652 | 0 | atomic_replace_range_.has_value() && !atomic_replace_range_->unset()); |
653 | 0 | if (!status.ok()) { |
654 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
655 | 0 | "Failed to check ranges overlap with memtables: %s", |
656 | 0 | status.ToString().c_str()); |
657 | 0 | } |
658 | 0 | } |
659 | 0 | if (status.ok() && *flush_needed) { |
660 | 0 | if (!ingestion_options_.allow_blocking_flush) { |
661 | 0 | status = Status::InvalidArgument("External file requires flush"); |
662 | 0 | } |
663 | 0 | if (ucmp_->timestamp_size() > 0) { |
664 | 0 | status = Status::InvalidArgument( |
665 | 0 | "Column family enables user-defined timestamps, please make " |
666 | 0 | "sure the key range (without timestamp) of external file does not " |
667 | 0 | "overlap with key range in the memtables."); |
668 | 0 | } |
669 | 0 | } |
670 | 0 | return status; |
671 | 0 | } |
672 | | |
673 | | // REQUIRES: we have become the only writer by entering both write_thread_ and |
674 | | // nonmem_write_thread_ |
675 | 0 | Status ExternalSstFileIngestionJob::Run() { |
676 | 0 | SuperVersion* super_version = cfd_->GetSuperVersion(); |
677 | | // If column family is flushed after Prepare and before Run, we should have a |
678 | | // specific state of Memtables. The mutable Memtable should be empty, and the |
679 | | // immutable Memtable list should be empty. |
680 | 0 | if (flushed_before_run_ && (super_version->imm->NumNotFlushed() != 0 || |
681 | 0 | !super_version->mem->IsEmpty())) { |
682 | 0 | return Status::TryAgain( |
683 | 0 | "Inconsistent memtable state detected when flushed before run."); |
684 | 0 | } |
685 | 0 | Status status; |
686 | | #ifndef NDEBUG |
687 | | // We should never run the job with a memtable that is overlapping |
688 | | // with the files we are ingesting |
689 | | bool need_flush = false; |
690 | | status = NeedsFlush(&need_flush, super_version); |
691 | | if (!status.ok()) { |
692 | | ROCKS_LOG_WARN(db_options_.info_log, |
693 | | "Failed to check if flush is needed: %s", |
694 | | status.ToString().c_str()); |
695 | | return status; |
696 | | } |
697 | | if (need_flush) { |
698 | | return Status::TryAgain("need_flush"); |
699 | | } |
700 | | assert(status.ok() && need_flush == false); |
701 | | #endif |
702 | |
|
703 | 0 | bool force_global_seqno = false; |
704 | |
|
705 | 0 | if (ingestion_options_.snapshot_consistency && !db_snapshots_->empty()) { |
706 | | // We need to assign a global sequence number to all the files even |
707 | | // if the don't overlap with any ranges since we have snapshots |
708 | 0 | force_global_seqno = true; |
709 | 0 | } |
710 | | // It is safe to use this instead of LastAllocatedSequence since we are |
711 | | // the only active writer, and hence they are equal |
712 | 0 | SequenceNumber last_seqno = versions_->LastSequence(); |
713 | 0 | edit_.SetColumnFamily(cfd_->GetID()); |
714 | |
|
715 | 0 | if (atomic_replace_range_.has_value()) { |
716 | 0 | auto* vstorage = super_version->current->storage_info(); |
717 | 0 | if (atomic_replace_range_->unset()) { |
718 | 0 | if (cfd_->compaction_picker()->IsCompactionInProgress()) { |
719 | 0 | return Status::InvalidArgument( |
720 | 0 | "Atomic replace range (full) overlaps with pending compaction"); |
721 | 0 | } |
722 | 0 | for (int lvl = 0; lvl < cfd_->NumberLevels(); lvl++) { |
723 | 0 | for (auto file : vstorage->LevelFiles(lvl)) { |
724 | | // Set up to delete file to be replaced |
725 | 0 | edit_.DeleteFile(lvl, file->fd.GetNumber()); |
726 | 0 | } |
727 | 0 | } |
728 | 0 | } else { |
729 | 0 | assert(!atomic_replace_range_->smallest_internal_key.unset()); |
730 | 0 | assert(!atomic_replace_range_->largest_internal_key.unset()); |
731 | 0 | bool has_partial_overlap = false; |
732 | 0 | for (int lvl = 0; lvl < cfd_->NumberLevels(); lvl++) { |
733 | 0 | if (cfd_->RangeOverlapWithCompaction( |
734 | 0 | atomic_replace_range_->smallest_internal_key.user_key(), |
735 | 0 | atomic_replace_range_->largest_internal_key.user_key(), lvl, |
736 | 0 | /*range_limit_exclusive=*/true)) { |
737 | 0 | return Status::InvalidArgument( |
738 | 0 | "Atomic replace range overlaps with pending compaction"); |
739 | 0 | } |
740 | 0 | for (auto file : vstorage->LevelFiles(lvl)) { |
741 | 0 | if (file_range_checker_.Overlaps(*atomic_replace_range_, |
742 | 0 | file->smallest, file->largest)) { |
743 | 0 | if (file_range_checker_.Contains(*atomic_replace_range_, |
744 | 0 | file->smallest, file->largest)) { |
745 | | // Set up to delete file to be replaced |
746 | 0 | edit_.DeleteFile(lvl, file->fd.GetNumber()); |
747 | 0 | } else { |
748 | 0 | has_partial_overlap = true; |
749 | 0 | } |
750 | 0 | } |
751 | 0 | } |
752 | 0 | } |
753 | 0 | if (has_partial_overlap) { |
754 | 0 | if (ingestion_options_.fail_if_not_bottommost_level && |
755 | 0 | SupportsAtomicReplaceRangeTombstone()) { |
756 | 0 | return Status::TryAgain( |
757 | 0 | "Atomic replace range partially overlaps with an existing file, " |
758 | 0 | "so replacement files cannot all be ingested to Lmax"); |
759 | 0 | } |
760 | 0 | if (!atomic_replace_range_tombstone_.has_value()) { |
761 | 0 | if (SupportsAtomicReplaceRangeTombstone()) { |
762 | 0 | return Status::TryAgain( |
763 | 0 | "Atomic replace range acquired a partial overlap after " |
764 | 0 | "preparation; retry ingestion"); |
765 | 0 | } |
766 | 0 | return Status::InvalidArgument( |
767 | 0 | "Atomic replace range partially overlaps with existing file"); |
768 | 0 | } |
769 | 0 | ActivateAtomicReplaceRangeTombstone(); |
770 | 0 | } |
771 | 0 | } |
772 | 0 | } |
773 | | |
774 | | // Find levels to ingest into |
775 | 0 | std::optional<int> prev_batch_uppermost_level; |
776 | | // Batches at the front contain older updates and are placed deeper in the |
777 | | // LSM tree than later overlapping batches. |
778 | 0 | for (auto& batch : file_batches_to_ingest_) { |
779 | 0 | int batch_uppermost_level = 0; |
780 | 0 | status = AssignLevelsForOneBatch(batch, super_version, force_global_seqno, |
781 | 0 | &last_seqno, &batch_uppermost_level, |
782 | 0 | prev_batch_uppermost_level); |
783 | 0 | if (!status.ok()) { |
784 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
785 | 0 | "Failed to assign levels for one batch: %s", |
786 | 0 | status.ToString().c_str()); |
787 | 0 | return status; |
788 | 0 | } |
789 | | |
790 | 0 | prev_batch_uppermost_level = batch_uppermost_level; |
791 | 0 | } |
792 | | |
793 | 0 | CreateEquivalentFileIngestingCompactions(); |
794 | 0 | return status; |
795 | 0 | } |
796 | | |
797 | | Status ExternalSstFileIngestionJob::AssignLevelsForOneBatch( |
798 | | FileBatchInfo& batch, SuperVersion* super_version, bool force_global_seqno, |
799 | | SequenceNumber* last_seqno, int* batch_uppermost_level, |
800 | 0 | std::optional<int> prev_batch_uppermost_level) { |
801 | 0 | Status status; |
802 | 0 | assert(batch_uppermost_level); |
803 | 0 | *batch_uppermost_level = std::numeric_limits<int>::max(); |
804 | 0 | for (IngestedFileInfo* file : batch.files) { |
805 | 0 | assert(file); |
806 | 0 | SequenceNumber assigned_seqno = 0; |
807 | 0 | if (ingestion_options_.ingest_behind) { |
808 | 0 | status = CheckLevelForIngestedBehindFile(file); |
809 | 0 | } else { |
810 | 0 | status = AssignLevelAndSeqnoForIngestedFile( |
811 | 0 | super_version, force_global_seqno, cfd_->ioptions().compaction_style, |
812 | 0 | *last_seqno, file, &assigned_seqno, prev_batch_uppermost_level); |
813 | 0 | } |
814 | | |
815 | | // Modify the smallest/largest internal key to include the sequence number |
816 | | // that we just learned. Only overwrite sequence number zero. There could |
817 | | // be a nonzero sequence number already to indicate a range tombstone's |
818 | | // exclusive endpoint. |
819 | 0 | ParsedInternalKey smallest_parsed, largest_parsed; |
820 | 0 | if (status.ok()) { |
821 | 0 | status = ParseInternalKey(*(file->smallest_internal_key.rep()), |
822 | 0 | &smallest_parsed, false /* log_err_key */); |
823 | 0 | } |
824 | 0 | if (status.ok()) { |
825 | 0 | status = ParseInternalKey(*(file->largest_internal_key.rep()), |
826 | 0 | &largest_parsed, false /* log_err_key */); |
827 | 0 | } |
828 | 0 | if (!status.ok()) { |
829 | 0 | ROCKS_LOG_WARN(db_options_.info_log, "Failed to parse internal key: %s", |
830 | 0 | status.ToString().c_str()); |
831 | 0 | return status; |
832 | 0 | } |
833 | | |
834 | | // If any ingested file overlaps with the DB, it will fail here. |
835 | 0 | if (ingestion_options_.allow_db_generated_files && assigned_seqno != 0) { |
836 | 0 | return Status::InvalidArgument( |
837 | 0 | "An ingested file overlaps with existing data in the DB and has been " |
838 | 0 | "assigned a non-zero sequence number, which is not allowed when " |
839 | 0 | "'allow_db_generated_files' is enabled."); |
840 | 0 | } |
841 | | |
842 | 0 | if (smallest_parsed.sequence == 0 && assigned_seqno != 0) { |
843 | 0 | UpdateInternalKey(file->smallest_internal_key.rep(), assigned_seqno, |
844 | 0 | smallest_parsed.type); |
845 | 0 | } |
846 | 0 | if (largest_parsed.sequence == 0 && assigned_seqno != 0) { |
847 | 0 | UpdateInternalKey(file->largest_internal_key.rep(), assigned_seqno, |
848 | 0 | largest_parsed.type); |
849 | 0 | } |
850 | |
|
851 | 0 | status = AssignGlobalSeqnoForIngestedFile(file, assigned_seqno); |
852 | 0 | if (!status.ok()) { |
853 | 0 | ROCKS_LOG_WARN( |
854 | 0 | db_options_.info_log, |
855 | 0 | "Failed to assign global sequence number for ingested file: %s", |
856 | 0 | status.ToString().c_str()); |
857 | 0 | return status; |
858 | 0 | } |
859 | 0 | TEST_SYNC_POINT_CALLBACK("ExternalSstFileIngestionJob::Run", |
860 | 0 | &assigned_seqno); |
861 | 0 | assert(assigned_seqno == 0 || assigned_seqno == *last_seqno + 1); |
862 | 0 | if (assigned_seqno > *last_seqno) { |
863 | 0 | *last_seqno = assigned_seqno; |
864 | 0 | } |
865 | 0 | max_assigned_seqno_ = std::max(max_assigned_seqno_, assigned_seqno); |
866 | |
|
867 | 0 | status = GenerateChecksumForIngestedFile(file); |
868 | 0 | if (!status.ok()) { |
869 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
870 | 0 | "Failed to generate checksum for ingested file: %s", |
871 | 0 | status.ToString().c_str()); |
872 | 0 | return status; |
873 | 0 | } |
874 | | |
875 | | // We use the import time as the ancester time. This is the time the data |
876 | | // is written to the database. |
877 | 0 | int64_t temp_current_time = 0; |
878 | 0 | uint64_t current_time = kUnknownFileCreationTime; |
879 | 0 | uint64_t oldest_ancester_time = kUnknownOldestAncesterTime; |
880 | 0 | if (clock_->GetCurrentTime(&temp_current_time).ok()) { |
881 | 0 | current_time = oldest_ancester_time = |
882 | 0 | static_cast<uint64_t>(temp_current_time); |
883 | 0 | } |
884 | 0 | uint64_t tail_size = FileMetaData::CalculateTailSize( |
885 | 0 | file->fd.GetFileSize(), file->table_properties); |
886 | |
|
887 | 0 | bool marked_for_compaction = |
888 | 0 | file->table_properties.num_range_deletions == 1 && |
889 | 0 | (file->table_properties.num_entries == |
890 | 0 | file->table_properties.num_range_deletions); |
891 | 0 | SequenceNumber smallest_seqno = file->assigned_seqno; |
892 | 0 | SequenceNumber largest_seqno = file->assigned_seqno; |
893 | 0 | if (ingestion_options_.allow_db_generated_files) { |
894 | 0 | assert(file->assigned_seqno == 0); |
895 | 0 | assert(file->smallest_seqno != kMaxSequenceNumber); |
896 | 0 | assert(file->largest_seqno != kMaxSequenceNumber); |
897 | 0 | smallest_seqno = file->smallest_seqno; |
898 | 0 | largest_seqno = file->largest_seqno; |
899 | 0 | max_assigned_seqno_ = std::max(max_assigned_seqno_, file->largest_seqno); |
900 | 0 | } |
901 | 0 | FileMetaData f_metadata( |
902 | 0 | file->fd.GetNumber(), file->fd.GetPathId(), file->fd.GetFileSize(), |
903 | 0 | file->smallest_internal_key, file->largest_internal_key, smallest_seqno, |
904 | 0 | largest_seqno, false, file->file_temperature, kInvalidBlobFileNumber, |
905 | 0 | oldest_ancester_time, current_time, |
906 | 0 | ingestion_options_.ingest_behind |
907 | 0 | ? kReservedEpochNumberForFileIngestedBehind |
908 | 0 | : cfd_->NewEpochNumber(), // orders files ingested to L0 |
909 | 0 | file->file_checksum, file->file_checksum_func_name, file->unique_id, 0, |
910 | 0 | tail_size, file->user_defined_timestamps_persisted, "", ""); |
911 | 0 | f_metadata.temperature = file->file_temperature; |
912 | 0 | f_metadata.marked_for_compaction = marked_for_compaction; |
913 | 0 | f_metadata.skip_index_and_filter_blocks_prefetch = |
914 | 0 | !file->prefetch_lmax_index_and_filter_blocks && |
915 | 0 | file->picked_level == cfd_->NumberLevels() - 1; |
916 | | // Extract min/max timestamps from table properties for UDT support. |
917 | | // This ensures ingested files have proper timestamp ranges in FileMetaData, |
918 | | // similar to files created by flush and compaction. |
919 | 0 | ExtractTimestampFromTableProperties(file->table_properties, &f_metadata); |
920 | | // Retrieve file open metadata for fast SST open |
921 | 0 | if (mutable_db_options_.fast_sst_open) { |
922 | 0 | std::unique_ptr<FSRandomAccessFile> readable_file; |
923 | 0 | FileOptions fopts{env_options_}; |
924 | 0 | fopts.file_checksum = f_metadata.file_checksum; |
925 | 0 | fopts.file_checksum_func_name = f_metadata.file_checksum_func_name; |
926 | 0 | IOStatus io_s = fs_->NewRandomAccessFile(file->internal_file_path, fopts, |
927 | 0 | &readable_file, nullptr); |
928 | 0 | if (io_s.ok()) { |
929 | 0 | io_s = |
930 | 0 | readable_file->GetFileOpenMetadata(&f_metadata.file_open_metadata); |
931 | 0 | if (io_s.ok() && !f_metadata.file_open_metadata.empty() && |
932 | 0 | f_metadata.file_open_metadata.size() <= |
933 | 0 | FSRandomAccessFile::kMaxFileOpenMetadataSize) { |
934 | 0 | RecordTick(db_options_.stats, FILE_OPEN_METADATA_RETRIEVED); |
935 | 0 | } else { |
936 | 0 | if (io_s.ok() && f_metadata.file_open_metadata.size() > |
937 | 0 | FSRandomAccessFile::kMaxFileOpenMetadataSize) { |
938 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
939 | 0 | "File open metadata for %s too large (%zu bytes), " |
940 | 0 | "ignoring", |
941 | 0 | file->internal_file_path.c_str(), |
942 | 0 | f_metadata.file_open_metadata.size()); |
943 | 0 | } |
944 | 0 | f_metadata.file_open_metadata.clear(); |
945 | 0 | } |
946 | 0 | } |
947 | 0 | } |
948 | 0 | edit_.AddFile(file->picked_level, f_metadata); |
949 | |
|
950 | 0 | *batch_uppermost_level = |
951 | 0 | std::min(*batch_uppermost_level, file->picked_level); |
952 | 0 | } |
953 | | |
954 | 0 | return Status::OK(); |
955 | 0 | } |
956 | | |
957 | 0 | void ExternalSstFileIngestionJob::CreateEquivalentFileIngestingCompactions() { |
958 | | // A map from output level to input of compactions equivalent to this |
959 | | // ingestion job. |
960 | | // TODO: simplify below logic to creating compaction per ingested file |
961 | | // instead of per output level, once we figure out how to treat ingested files |
962 | | // with adjacent range deletion tombstones to same output level in the same |
963 | | // job as non-overlapping compactions. |
964 | 0 | std::map<int, CompactionInputFiles> |
965 | 0 | output_level_to_file_ingesting_compaction_input; |
966 | |
|
967 | 0 | for (const auto& pair : edit_.GetNewFiles()) { |
968 | 0 | int output_level = pair.first; |
969 | 0 | const FileMetaData& f_metadata = pair.second; |
970 | |
|
971 | 0 | CompactionInputFiles& input = |
972 | 0 | output_level_to_file_ingesting_compaction_input[output_level]; |
973 | 0 | if (input.files.empty()) { |
974 | | // Treat the source level of ingested files to be level 0 |
975 | 0 | input.level = 0; |
976 | 0 | } |
977 | |
|
978 | 0 | compaction_input_metdatas_.push_back(new FileMetaData(f_metadata)); |
979 | 0 | input.files.push_back(compaction_input_metdatas_.back()); |
980 | 0 | } |
981 | |
|
982 | 0 | for (const auto& pair : output_level_to_file_ingesting_compaction_input) { |
983 | 0 | int output_level = pair.first; |
984 | 0 | const CompactionInputFiles& input = pair.second; |
985 | |
|
986 | 0 | const auto& mutable_cf_options = cfd_->GetLatestMutableCFOptions(); |
987 | 0 | file_ingesting_compactions_.push_back(new Compaction( |
988 | 0 | cfd_->current()->storage_info(), cfd_->ioptions(), mutable_cf_options, |
989 | 0 | mutable_db_options_, {input}, output_level, |
990 | | /* output file size limit not applicable */ |
991 | 0 | MaxFileSizeForLevel(mutable_cf_options, output_level, |
992 | 0 | cfd_->ioptions().compaction_style), |
993 | 0 | LLONG_MAX /* max compaction bytes, not applicable */, |
994 | 0 | 0 /* output path ID, not applicable */, mutable_cf_options.compression, |
995 | 0 | mutable_cf_options.compression_opts, Temperature::kUnknown, |
996 | 0 | 0 /* max_subcompaction, not applicable */, |
997 | 0 | {} /* grandparents, not applicable */, |
998 | 0 | std::nullopt /* earliest_snapshot */, nullptr /* snapshot_checker */, |
999 | 0 | CompactionReason::kExternalSstIngestion, "" /* trim_ts */, |
1000 | 0 | -1 /* score, not applicable */, |
1001 | 0 | files_overlap_ /* l0_files_might_overlap, not applicable */)); |
1002 | 0 | } |
1003 | 0 | } |
1004 | | |
1005 | 0 | void ExternalSstFileIngestionJob::RegisterRange() { |
1006 | 0 | for (const auto& c : file_ingesting_compactions_) { |
1007 | 0 | cfd_->compaction_picker()->RegisterCompaction(c); |
1008 | 0 | } |
1009 | 0 | } |
1010 | | |
1011 | 0 | void ExternalSstFileIngestionJob::UnregisterRange() { |
1012 | 0 | for (const auto& c : file_ingesting_compactions_) { |
1013 | 0 | cfd_->compaction_picker()->UnregisterCompaction(c); |
1014 | 0 | delete c; |
1015 | 0 | } |
1016 | 0 | file_ingesting_compactions_.clear(); |
1017 | |
|
1018 | 0 | for (const auto& f : compaction_input_metdatas_) { |
1019 | 0 | delete f; |
1020 | 0 | } |
1021 | 0 | compaction_input_metdatas_.clear(); |
1022 | 0 | } |
1023 | | |
1024 | 0 | void ExternalSstFileIngestionJob::UpdateStats() { |
1025 | | // Update internal stats for new ingested files |
1026 | 0 | uint64_t total_keys = 0; |
1027 | 0 | uint64_t total_l0_files = 0; |
1028 | 0 | uint64_t total_time = clock_->NowMicros() - job_start_time_; |
1029 | |
|
1030 | 0 | EventLoggerStream stream = event_logger_->Log(); |
1031 | 0 | stream << "event" << "ingest_finished"; |
1032 | 0 | stream << "files_ingested"; |
1033 | 0 | stream.StartArray(); |
1034 | |
|
1035 | 0 | for (IngestedFileInfo& f : files_to_ingest_) { |
1036 | 0 | InternalStats::CompactionStats stats( |
1037 | 0 | CompactionReason::kExternalSstIngestion, 1); |
1038 | 0 | stats.micros = total_time; |
1039 | | // If actual copy occurred for this file, then we need to count the file |
1040 | | // size as the actual bytes written. If the file was linked, then we ignore |
1041 | | // the bytes written for file metadata. |
1042 | | // TODO (yanqin) maybe account for file metadata bytes for exact accuracy? |
1043 | 0 | if (f.copy_file) { |
1044 | 0 | stats.bytes_written = f.fd.GetFileSize(); |
1045 | 0 | } else { |
1046 | 0 | stats.bytes_moved = f.fd.GetFileSize(); |
1047 | 0 | } |
1048 | 0 | stats.num_output_files = 1; |
1049 | 0 | cfd_->internal_stats()->AddCompactionStats(f.picked_level, |
1050 | 0 | Env::Priority::USER, stats); |
1051 | 0 | cfd_->internal_stats()->AddCFStats(InternalStats::BYTES_INGESTED_ADD_FILE, |
1052 | 0 | f.fd.GetFileSize()); |
1053 | 0 | total_keys += f.num_entries; |
1054 | 0 | if (f.picked_level == 0) { |
1055 | 0 | total_l0_files += 1; |
1056 | 0 | } |
1057 | 0 | ROCKS_LOG_INFO( |
1058 | 0 | db_options_.info_log, |
1059 | 0 | "[AddFile] External SST file %s was ingested in L%d with path %s " |
1060 | 0 | "(global_seqno=%" PRIu64 ")\n", |
1061 | 0 | f.external_file_path.c_str(), f.picked_level, |
1062 | 0 | f.internal_file_path.c_str(), f.assigned_seqno); |
1063 | 0 | stream << "file" << f.internal_file_path << "level" << f.picked_level; |
1064 | 0 | } |
1065 | 0 | stream.EndArray(); |
1066 | |
|
1067 | 0 | stream << "lsm_state"; |
1068 | 0 | stream.StartArray(); |
1069 | 0 | auto vstorage = cfd_->current()->storage_info(); |
1070 | 0 | for (int level = 0; level < vstorage->num_levels(); ++level) { |
1071 | 0 | stream << vstorage->NumLevelFiles(level); |
1072 | 0 | } |
1073 | 0 | stream.EndArray(); |
1074 | |
|
1075 | 0 | cfd_->internal_stats()->AddCFStats(InternalStats::INGESTED_NUM_KEYS_TOTAL, |
1076 | 0 | total_keys); |
1077 | 0 | cfd_->internal_stats()->AddCFStats(InternalStats::INGESTED_NUM_FILES_TOTAL, |
1078 | 0 | files_to_ingest_.size()); |
1079 | 0 | cfd_->internal_stats()->AddCFStats( |
1080 | 0 | InternalStats::INGESTED_LEVEL0_NUM_FILES_TOTAL, total_l0_files); |
1081 | 0 | } |
1082 | | |
1083 | 0 | void ExternalSstFileIngestionJob::Cleanup(const Status& status) { |
1084 | 0 | IOOptions io_opts; |
1085 | 0 | if (!status.ok()) { |
1086 | | // We failed to add the files to the database |
1087 | | // remove all the files we copied |
1088 | 0 | DeleteInternalFiles(); |
1089 | 0 | files_overlap_ = false; |
1090 | 0 | } else { |
1091 | 0 | if (atomic_replace_range_tombstone_.has_value()) { |
1092 | 0 | Status s = |
1093 | 0 | fs_->DeleteFile(atomic_replace_range_tombstone_->internal_file_path, |
1094 | 0 | io_opts, nullptr); |
1095 | 0 | if (!s.ok()) { |
1096 | 0 | ROCKS_LOG_WARN( |
1097 | 0 | db_options_.info_log, |
1098 | 0 | "Failed to remove unused atomic replace range tombstone %s: %s", |
1099 | 0 | atomic_replace_range_tombstone_->internal_file_path.c_str(), |
1100 | 0 | s.ToString().c_str()); |
1101 | 0 | } |
1102 | 0 | atomic_replace_range_tombstone_.reset(); |
1103 | 0 | } |
1104 | 0 | } |
1105 | 0 | if (status.ok() && ingestion_options_.move_files) { |
1106 | | // The files were moved and added successfully, remove original file links |
1107 | 0 | for (IngestedFileInfo& f : files_to_ingest_) { |
1108 | 0 | if (f.generated_for_ingestion) { |
1109 | 0 | continue; |
1110 | 0 | } |
1111 | 0 | Status s = fs_->DeleteFile(f.external_file_path, io_opts, nullptr); |
1112 | 0 | if (!s.ok()) { |
1113 | 0 | ROCKS_LOG_WARN( |
1114 | 0 | db_options_.info_log, |
1115 | 0 | "%s was added to DB successfully but failed to remove original " |
1116 | 0 | "file link : %s", |
1117 | 0 | f.external_file_path.c_str(), s.ToString().c_str()); |
1118 | 0 | } |
1119 | 0 | } |
1120 | 0 | } |
1121 | 0 | } |
1122 | | |
1123 | 0 | void ExternalSstFileIngestionJob::DeleteInternalFiles() { |
1124 | 0 | IOOptions io_opts; |
1125 | 0 | for (IngestedFileInfo& f : files_to_ingest_) { |
1126 | 0 | if (f.internal_file_path.empty()) { |
1127 | 0 | continue; |
1128 | 0 | } |
1129 | 0 | Status s = fs_->DeleteFile(f.internal_file_path, io_opts, nullptr); |
1130 | 0 | if (!s.ok()) { |
1131 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
1132 | 0 | "AddFile() clean up for file %s failed : %s", |
1133 | 0 | f.internal_file_path.c_str(), s.ToString().c_str()); |
1134 | 0 | } |
1135 | 0 | } |
1136 | 0 | if (atomic_replace_range_tombstone_.has_value() && |
1137 | 0 | !atomic_replace_range_tombstone_->internal_file_path.empty()) { |
1138 | 0 | Status s = fs_->DeleteFile( |
1139 | 0 | atomic_replace_range_tombstone_->internal_file_path, io_opts, nullptr); |
1140 | 0 | if (!s.ok()) { |
1141 | 0 | ROCKS_LOG_WARN( |
1142 | 0 | db_options_.info_log, |
1143 | 0 | "AddFile() clean up for generated tombstone file %s failed : %s", |
1144 | 0 | atomic_replace_range_tombstone_->internal_file_path.c_str(), |
1145 | 0 | s.ToString().c_str()); |
1146 | 0 | } |
1147 | 0 | atomic_replace_range_tombstone_.reset(); |
1148 | 0 | } |
1149 | 0 | } |
1150 | | |
1151 | | Status ExternalSstFileIngestionJob::ResetTableReader( |
1152 | | const std::string& external_file, uint64_t new_file_number, |
1153 | | bool user_defined_timestamps_persisted, SuperVersion* sv, |
1154 | | IngestedFileInfo* file_to_ingest, |
1155 | 0 | std::unique_ptr<TableReader>* table_reader) { |
1156 | 0 | std::unique_ptr<FSRandomAccessFile> sst_file; |
1157 | 0 | FileOptions fo{env_options_}; |
1158 | 0 | fo.temperature = file_to_ingest->file_temperature; |
1159 | 0 | fo.file_checksum_func_name = kNoFileChecksumFuncName; |
1160 | 0 | Status status = |
1161 | 0 | fs_->NewRandomAccessFile(external_file, fo, &sst_file, nullptr); |
1162 | 0 | if (!status.ok()) { |
1163 | 0 | ROCKS_LOG_WARN( |
1164 | 0 | db_options_.info_log, |
1165 | 0 | "Failed to create random access file for external file %s: %s", |
1166 | 0 | external_file.c_str(), status.ToString().c_str()); |
1167 | 0 | return status; |
1168 | 0 | } |
1169 | 0 | Temperature updated_temp = sst_file->GetTemperature(); |
1170 | 0 | if (updated_temp != Temperature::kUnknown && |
1171 | 0 | updated_temp != file_to_ingest->file_temperature) { |
1172 | | // The hint was missing or wrong. Track temperature reported by storage. |
1173 | 0 | file_to_ingest->file_temperature = updated_temp; |
1174 | 0 | } |
1175 | 0 | std::unique_ptr<RandomAccessFileReader> sst_file_reader( |
1176 | 0 | new RandomAccessFileReader(std::move(sst_file), external_file, |
1177 | 0 | nullptr /*Env*/, io_tracer_)); |
1178 | 0 | table_reader->reset(); |
1179 | 0 | ReadOptions ro; |
1180 | 0 | ro.fill_cache = ingestion_options_.fill_cache; |
1181 | 0 | status = sv->mutable_cf_options.table_factory->NewTableReader( |
1182 | 0 | ro, |
1183 | 0 | TableReaderOptions( |
1184 | 0 | cfd_->ioptions(), sv->mutable_cf_options.prefix_extractor, |
1185 | 0 | sv->mutable_cf_options.compression_manager.get(), env_options_, |
1186 | 0 | cfd_->internal_comparator(), |
1187 | 0 | sv->mutable_cf_options.block_protection_bytes_per_key, |
1188 | 0 | /*skip_filters*/ false, /*immortal*/ false, |
1189 | 0 | /*force_direct_prefetch*/ false, /*level*/ -1, |
1190 | 0 | /*block_cache_tracer*/ nullptr, |
1191 | 0 | /*max_file_size_for_l0_meta_pin*/ 0, versions_->DbSessionId(), |
1192 | 0 | /*cur_file_num*/ new_file_number, |
1193 | 0 | /* unique_id */ {}, /* largest_seqno */ 0, |
1194 | 0 | /* tail_size */ 0, user_defined_timestamps_persisted), |
1195 | 0 | std::move(sst_file_reader), file_to_ingest->file_size, table_reader, |
1196 | | // No need to prefetch index/filter if caching is not needed. |
1197 | 0 | /*prefetch_index_and_filter_in_cache=*/ingestion_options_.fill_cache); |
1198 | 0 | return status; |
1199 | 0 | } |
1200 | | |
1201 | | Status ExternalSstFileIngestionJob::SanityCheckTableProperties( |
1202 | | const std::string& external_file, const TableProperties& props, |
1203 | 0 | IngestedFileInfo* file_to_ingest) { |
1204 | 0 | const auto& uprops = props.user_collected_properties; |
1205 | | |
1206 | | // Get table version |
1207 | 0 | auto version_iter = uprops.find(ExternalSstFilePropertyNames::kVersion); |
1208 | 0 | if (version_iter == uprops.end()) { |
1209 | 0 | assert(!SstFileWriter::CreatedBySstFileWriter(props)); |
1210 | 0 | if (!ingestion_options_.allow_db_generated_files) { |
1211 | 0 | return Status::Corruption("External file version not found"); |
1212 | 0 | } else { |
1213 | | // 0 is special version for when a file from live DB does not have the |
1214 | | // version table property |
1215 | 0 | file_to_ingest->version = 0; |
1216 | 0 | } |
1217 | 0 | } else { |
1218 | 0 | assert(SstFileWriter::CreatedBySstFileWriter(props)); |
1219 | 0 | file_to_ingest->version = DecodeFixed32(version_iter->second.c_str()); |
1220 | 0 | } |
1221 | | |
1222 | 0 | auto seqno_iter = uprops.find(ExternalSstFilePropertyNames::kGlobalSeqno); |
1223 | 0 | if (file_to_ingest->version == 2) { |
1224 | | // version 2 imply that we have global sequence number |
1225 | 0 | if (seqno_iter == uprops.end()) { |
1226 | 0 | return Status::Corruption( |
1227 | 0 | "External file global sequence number not found"); |
1228 | 0 | } |
1229 | | |
1230 | | // Set the global sequence number |
1231 | 0 | file_to_ingest->original_seqno = DecodeFixed64(seqno_iter->second.c_str()); |
1232 | 0 | file_to_ingest->global_seqno_offset = |
1233 | 0 | static_cast<size_t>(props.external_sst_file_global_seqno_offset); |
1234 | | // The on-disk offset is only needed if we will write the global seqno back |
1235 | | // into the file (write_global_seqno). The metadata fast-path does not open |
1236 | | // the file, and its in-memory table properties do not carry the offset |
1237 | | // (it is only computed while reading the file back); that is fine as long |
1238 | | // as write_global_seqno is not requested. |
1239 | 0 | if (ingestion_options_.write_global_seqno && |
1240 | 0 | file_to_ingest->global_seqno_offset == 0) { |
1241 | 0 | return Status::Corruption("Was not able to find file global seqno field"); |
1242 | 0 | } |
1243 | 0 | } else if (file_to_ingest->version == 1) { |
1244 | | // SST file V1 should not have global seqno field |
1245 | 0 | assert(seqno_iter == uprops.end()); |
1246 | 0 | file_to_ingest->original_seqno = 0; |
1247 | 0 | if (ingestion_options_.allow_blocking_flush || |
1248 | 0 | ingestion_options_.allow_global_seqno) { |
1249 | 0 | return Status::InvalidArgument( |
1250 | 0 | "External SST file V1 does not support global seqno"); |
1251 | 0 | } |
1252 | 0 | } else if (file_to_ingest->version == 0) { |
1253 | | // allow_db_generated_files is true |
1254 | 0 | assert(seqno_iter == uprops.end()); |
1255 | 0 | file_to_ingest->original_seqno = 0; |
1256 | 0 | file_to_ingest->global_seqno_offset = 0; |
1257 | 0 | } else { |
1258 | 0 | return Status::InvalidArgument("External file version " + |
1259 | 0 | std::to_string(file_to_ingest->version) + |
1260 | 0 | " is not supported"); |
1261 | 0 | } |
1262 | | |
1263 | 0 | file_to_ingest->cf_id = static_cast<uint32_t>(props.column_family_id); |
1264 | 0 | file_to_ingest->table_properties = props; |
1265 | | |
1266 | | // Get number of entries in table |
1267 | 0 | file_to_ingest->num_entries = props.num_entries; |
1268 | 0 | file_to_ingest->num_range_deletions = props.num_range_deletions; |
1269 | | |
1270 | | // Validate table properties related to comparator name and user defined |
1271 | | // timestamps persisted flag. |
1272 | 0 | file_to_ingest->user_defined_timestamps_persisted = |
1273 | 0 | static_cast<bool>(props.user_defined_timestamps_persisted); |
1274 | 0 | bool mark_sst_file_has_no_udt = false; |
1275 | 0 | Status s = ValidateUserDefinedTimestampsOptions( |
1276 | 0 | cfd_->user_comparator(), props.comparator_name, |
1277 | 0 | cfd_->ioptions().persist_user_defined_timestamps, |
1278 | 0 | file_to_ingest->user_defined_timestamps_persisted, |
1279 | 0 | &mark_sst_file_has_no_udt); |
1280 | 0 | if (s.ok() && mark_sst_file_has_no_udt) { |
1281 | | // A column family that enables user-defined timestamps in Memtable only |
1282 | | // feature can also ingest external files created by a setting that disables |
1283 | | // user-defined timestamps. In that case, we need to re-mark the |
1284 | | // user_defined_timestamps_persisted flag for the file. The open-and-scan |
1285 | | // caller is then responsible for reopening its `TableReader` with the |
1286 | | // updated flag. |
1287 | 0 | file_to_ingest->user_defined_timestamps_persisted = false; |
1288 | 0 | } else if (!s.ok()) { |
1289 | 0 | ROCKS_LOG_WARN( |
1290 | 0 | db_options_.info_log, |
1291 | 0 | "ValidateUserDefinedTimestampsOptions failed for external file %s: %s", |
1292 | 0 | external_file.c_str(), s.ToString().c_str()); |
1293 | 0 | return s; |
1294 | 0 | } |
1295 | | |
1296 | 0 | return s; |
1297 | 0 | } |
1298 | | |
1299 | | Status ExternalSstFileIngestionJob::GetIngestedFileInfoFromFileInfo( |
1300 | | const std::string& external_file, |
1301 | | const PreparedFileInfo& prepared_file_info, |
1302 | 0 | IngestedFileInfo* file_to_ingest) { |
1303 | | // Boundaries, size, and table properties were obtained without opening the |
1304 | | // file (e.g. produced by SstFileWriter::Finish). We reuse them directly. |
1305 | 0 | file_to_ingest->file_size = prepared_file_info.file_size; |
1306 | 0 | Status status = SanityCheckTableProperties( |
1307 | 0 | external_file, prepared_file_info.table_properties, file_to_ingest); |
1308 | 0 | if (!status.ok()) { |
1309 | 0 | ROCKS_LOG_WARN( |
1310 | 0 | db_options_.info_log, |
1311 | 0 | "Failed to sanity check table properties for external file %s: %s", |
1312 | 0 | external_file.c_str(), status.ToString().c_str()); |
1313 | 0 | return status; |
1314 | 0 | } |
1315 | | |
1316 | 0 | const size_t ts_sz = ucmp_->timestamp_size(); |
1317 | 0 | if (ts_sz > 0 && !file_to_ingest->user_defined_timestamps_persisted) { |
1318 | 0 | auto pad_timestamp = [ts_sz](std::string* result, const Slice& key) { |
1319 | 0 | assert(result->empty()); |
1320 | 0 | if (ExtractValueType(key) == kTypeRangeDeletion) { |
1321 | 0 | PadInternalKeyWithMaxTimestamp(result, key, ts_sz); |
1322 | 0 | } else { |
1323 | 0 | PadInternalKeyWithMinTimestamp(result, key, ts_sz); |
1324 | 0 | } |
1325 | 0 | }; |
1326 | 0 | pad_timestamp(file_to_ingest->smallest_internal_key.rep(), |
1327 | 0 | prepared_file_info.smallest.Encode()); |
1328 | 0 | pad_timestamp(file_to_ingest->largest_internal_key.rep(), |
1329 | 0 | prepared_file_info.largest.Encode()); |
1330 | 0 | } else { |
1331 | 0 | file_to_ingest->smallest_internal_key = prepared_file_info.smallest; |
1332 | 0 | file_to_ingest->largest_internal_key = prepared_file_info.largest; |
1333 | 0 | } |
1334 | |
|
1335 | 0 | if (ingestion_options_.allow_db_generated_files) { |
1336 | | // Sequence numbers are preserved (not reassigned), so the bounds must be |
1337 | | // known before we skip the GetSeqnoBoundaryForFile scan. |
1338 | 0 | if (!file_to_ingest->table_properties.HasKeyLargestSeqno()) { |
1339 | 0 | return Status::Corruption("Unknown largest seqno for db generated file."); |
1340 | 0 | } |
1341 | 0 | file_to_ingest->largest_seqno = |
1342 | 0 | file_to_ingest->table_properties.key_largest_seqno; |
1343 | 0 | if (file_to_ingest->largest_seqno == 0) { |
1344 | 0 | file_to_ingest->smallest_seqno = 0; |
1345 | 0 | } else { |
1346 | 0 | if (!file_to_ingest->table_properties.HasKeySmallestSeqno()) { |
1347 | 0 | return Status::Corruption( |
1348 | 0 | "Unknown smallest seqno for db generated file."); |
1349 | 0 | } |
1350 | 0 | file_to_ingest->smallest_seqno = |
1351 | 0 | file_to_ingest->table_properties.key_smallest_seqno; |
1352 | 0 | } |
1353 | 0 | } else { |
1354 | | // Normal ingestion reassigns a global sequence number later, so the file's |
1355 | | // keys must currently be at seqno 0 (mirror of the open-and-scan check). |
1356 | 0 | SequenceNumber largest_seqno = |
1357 | 0 | file_to_ingest->table_properties.key_largest_seqno; |
1358 | | // UINT64_MAX means unknown and the file is generated before table property |
1359 | | // `key_largest_seqno` is introduced. |
1360 | 0 | if (largest_seqno != UINT64_MAX && largest_seqno > 0) { |
1361 | 0 | return Status::Corruption( |
1362 | 0 | "External file has non zero largest sequence number " + |
1363 | 0 | std::to_string(largest_seqno)); |
1364 | 0 | } |
1365 | 0 | } |
1366 | 0 | return Status::OK(); |
1367 | 0 | } |
1368 | | |
1369 | | Status ExternalSstFileIngestionJob::GetIngestedFileInfoFromFile( |
1370 | | const std::string& external_file, uint64_t new_file_number, |
1371 | | IngestedFileInfo* file_to_ingest, SuperVersion* sv, |
1372 | 0 | std::unique_ptr<TableReader>* out_table_reader) { |
1373 | 0 | TEST_SYNC_POINT("ExternalSstFileIngestionJob::GetIngestedFileInfo:ReadPath"); |
1374 | | |
1375 | | // Get external file size |
1376 | 0 | Status status = fs_->GetFileSize(external_file, IOOptions(), |
1377 | 0 | &file_to_ingest->file_size, nullptr); |
1378 | 0 | if (!status.ok()) { |
1379 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
1380 | 0 | "Failed to get file size for external file %s: %s", |
1381 | 0 | external_file.c_str(), status.ToString().c_str()); |
1382 | 0 | return status; |
1383 | 0 | } |
1384 | | |
1385 | | // Create TableReader for external file. |
1386 | 0 | std::unique_ptr<TableReader> table_reader; |
1387 | | // Initially create the `TableReader` with flag |
1388 | | // `user_defined_timestamps_persisted` to be true since that's the most |
1389 | | // common case |
1390 | 0 | status = ResetTableReader(external_file, new_file_number, |
1391 | 0 | /*user_defined_timestamps_persisted=*/true, sv, |
1392 | 0 | file_to_ingest, &table_reader); |
1393 | 0 | if (!status.ok()) { |
1394 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
1395 | 0 | "Failed to reset table reader for external file %s: %s", |
1396 | 0 | external_file.c_str(), status.ToString().c_str()); |
1397 | 0 | return status; |
1398 | 0 | } |
1399 | | |
1400 | 0 | status = SanityCheckTableProperties( |
1401 | 0 | external_file, *table_reader->GetTableProperties(), file_to_ingest); |
1402 | 0 | if (!status.ok()) { |
1403 | 0 | ROCKS_LOG_WARN( |
1404 | 0 | db_options_.info_log, |
1405 | 0 | "Failed to sanity check table properties for external file %s: %s", |
1406 | 0 | external_file.c_str(), status.ToString().c_str()); |
1407 | 0 | return status; |
1408 | 0 | } |
1409 | | |
1410 | | // The `TableReader` above was opened with `user_defined_timestamps_persisted` |
1411 | | // assumed true. If the sanity check determined the file has no persisted |
1412 | | // timestamps (UDT-in-Memtable-only feature), reopen it with the corrected |
1413 | | // flag so keys are parsed properly by the scan below. |
1414 | 0 | if (ucmp_->timestamp_size() > 0 && |
1415 | 0 | !file_to_ingest->user_defined_timestamps_persisted) { |
1416 | 0 | status = ResetTableReader(external_file, new_file_number, |
1417 | 0 | file_to_ingest->user_defined_timestamps_persisted, |
1418 | 0 | sv, file_to_ingest, &table_reader); |
1419 | 0 | if (!status.ok()) { |
1420 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
1421 | 0 | "Failed to reset table reader for external file %s: %s", |
1422 | 0 | external_file.c_str(), status.ToString().c_str()); |
1423 | 0 | return status; |
1424 | 0 | } |
1425 | 0 | } |
1426 | | |
1427 | 0 | const bool allow_data_in_errors = db_options_.allow_data_in_errors; |
1428 | 0 | ParsedInternalKey key; |
1429 | 0 | if (ingestion_options_.allow_db_generated_files) { |
1430 | | // We are ingesting a DB generated SST file for which we don't reassign |
1431 | | // sequence numbers. We need its smallest sequence number and largest |
1432 | | // sequence number for FileMetaData. |
1433 | 0 | Status seqno_status = GetSeqnoBoundaryForFile( |
1434 | 0 | table_reader.get(), sv, file_to_ingest, allow_data_in_errors); |
1435 | |
|
1436 | 0 | if (!seqno_status.ok()) { |
1437 | 0 | ROCKS_LOG_WARN( |
1438 | 0 | db_options_.info_log, |
1439 | 0 | "Failed to get sequence number boundary for external file %s: %s", |
1440 | 0 | external_file.c_str(), seqno_status.ToString().c_str()); |
1441 | 0 | return seqno_status; |
1442 | 0 | } |
1443 | 0 | assert(file_to_ingest->smallest_seqno <= file_to_ingest->largest_seqno); |
1444 | 0 | assert(file_to_ingest->largest_seqno < kMaxSequenceNumber); |
1445 | 0 | } else { |
1446 | 0 | SequenceNumber largest_seqno = |
1447 | 0 | table_reader.get()->GetTableProperties()->key_largest_seqno; |
1448 | | // UINT64_MAX means unknown and the file is generated before table |
1449 | | // property `key_largest_seqno` is introduced. |
1450 | 0 | if (largest_seqno != UINT64_MAX && largest_seqno > 0) { |
1451 | 0 | return Status::Corruption( |
1452 | 0 | "External file has non zero largest sequence number " + |
1453 | 0 | std::to_string(largest_seqno)); |
1454 | 0 | } |
1455 | 0 | } |
1456 | | |
1457 | | // TODO: plumb Env::IOActivity, Env::IOPriority |
1458 | 0 | ReadOptions ro; |
1459 | 0 | ro.fill_cache = ingestion_options_.fill_cache; |
1460 | 0 | std::unique_ptr<InternalIterator> iter(table_reader->NewIterator( |
1461 | 0 | ro, sv->mutable_cf_options.prefix_extractor.get(), /*arena=*/nullptr, |
1462 | 0 | /*skip_filters=*/false, TableReaderCaller::kExternalSSTIngestion)); |
1463 | | |
1464 | | // Get first (smallest) and last (largest) key from file. |
1465 | 0 | iter->SeekToFirst(); |
1466 | 0 | if (iter->Valid()) { |
1467 | 0 | Status pik_status = |
1468 | 0 | ParseInternalKey(iter->key(), &key, allow_data_in_errors); |
1469 | 0 | if (!pik_status.ok()) { |
1470 | 0 | return Status::Corruption("Corrupted key in external file. ", |
1471 | 0 | pik_status.getState()); |
1472 | 0 | } |
1473 | 0 | if (key.sequence != 0 && !ingestion_options_.allow_db_generated_files) { |
1474 | 0 | return Status::Corruption("External file has non zero sequence number"); |
1475 | 0 | } |
1476 | 0 | file_to_ingest->smallest_internal_key.SetFrom(key); |
1477 | |
|
1478 | 0 | Slice largest; |
1479 | 0 | if (strcmp(sv->mutable_cf_options.table_factory->Name(), "PlainTable") == |
1480 | 0 | 0) { |
1481 | | // PlainTable iterator does not support SeekToLast(). |
1482 | 0 | largest = iter->key(); |
1483 | 0 | for (; iter->Valid(); iter->Next()) { |
1484 | 0 | if (cfd_->internal_comparator().Compare(iter->key(), largest) > 0) { |
1485 | 0 | largest = iter->key(); |
1486 | 0 | } |
1487 | 0 | } |
1488 | 0 | if (!iter->status().ok()) { |
1489 | 0 | return iter->status(); |
1490 | 0 | } |
1491 | 0 | } else { |
1492 | 0 | iter->SeekToLast(); |
1493 | 0 | if (!iter->Valid()) { |
1494 | 0 | if (iter->status().ok()) { |
1495 | | // The file contains at least 1 key since iter is valid after |
1496 | | // SeekToFirst(). |
1497 | 0 | return Status::Corruption("Can not find largest key in sst file"); |
1498 | 0 | } else { |
1499 | 0 | return iter->status(); |
1500 | 0 | } |
1501 | 0 | } |
1502 | 0 | largest = iter->key(); |
1503 | 0 | } |
1504 | | |
1505 | 0 | pik_status = ParseInternalKey(largest, &key, allow_data_in_errors); |
1506 | 0 | if (!pik_status.ok()) { |
1507 | 0 | return Status::Corruption("Corrupted key in external file. ", |
1508 | 0 | pik_status.getState()); |
1509 | 0 | } |
1510 | 0 | if (key.sequence != 0 && !ingestion_options_.allow_db_generated_files) { |
1511 | 0 | return Status::Corruption("External file has non zero sequence number"); |
1512 | 0 | } |
1513 | 0 | file_to_ingest->largest_internal_key.SetFrom(key); |
1514 | 0 | } else if (!iter->status().ok()) { |
1515 | 0 | return iter->status(); |
1516 | 0 | } |
1517 | | |
1518 | 0 | std::unique_ptr<InternalIterator> range_del_iter( |
1519 | 0 | table_reader->NewRangeTombstoneIterator(ro)); |
1520 | | // We may need to adjust these key bounds, depending on whether any range |
1521 | | // deletion tombstones extend past them. |
1522 | 0 | if (range_del_iter != nullptr) { |
1523 | 0 | for (range_del_iter->SeekToFirst(); range_del_iter->Valid(); |
1524 | 0 | range_del_iter->Next()) { |
1525 | 0 | Status pik_status = |
1526 | 0 | ParseInternalKey(range_del_iter->key(), &key, allow_data_in_errors); |
1527 | 0 | if (!pik_status.ok()) { |
1528 | 0 | return Status::Corruption("Corrupted key in external file. ", |
1529 | 0 | pik_status.getState()); |
1530 | 0 | } |
1531 | 0 | if (key.sequence != 0 && !ingestion_options_.allow_db_generated_files) { |
1532 | 0 | return Status::Corruption( |
1533 | 0 | "External file has a range deletion with non zero sequence " |
1534 | 0 | "number."); |
1535 | 0 | } |
1536 | | #ifndef NDEBUG |
1537 | | // To keep aligned with the fast path that expects range deletion keys to |
1538 | | // always use max ts. |
1539 | | const size_t ts_sz = ucmp_->timestamp_size(); |
1540 | | if (ts_sz > 0) { |
1541 | | const std::string max_ts(ts_sz, '\xff'); |
1542 | | assert(key.user_key.size() >= ts_sz); |
1543 | | assert(ucmp_->CompareTimestamp( |
1544 | | ExtractTimestampFromUserKey(key.user_key, ts_sz), max_ts) == |
1545 | | 0); |
1546 | | assert(range_del_iter->value().size() >= ts_sz); |
1547 | | assert(ucmp_->CompareTimestamp( |
1548 | | ExtractTimestampFromUserKey(range_del_iter->value(), ts_sz), |
1549 | | max_ts) == 0); |
1550 | | } |
1551 | | #endif |
1552 | 0 | RangeTombstone tombstone(key, range_del_iter->value()); |
1553 | 0 | file_range_checker_.MaybeUpdateRange(tombstone.SerializeKey(), |
1554 | 0 | tombstone.SerializeEndKey(), |
1555 | 0 | file_to_ingest); |
1556 | 0 | } |
1557 | 0 | } |
1558 | | |
1559 | 0 | *out_table_reader = std::move(table_reader); |
1560 | 0 | return Status::OK(); |
1561 | 0 | } |
1562 | | |
1563 | | Status ExternalSstFileIngestionJob::GetIngestedFileInfo( |
1564 | | const std::string& external_file, uint64_t new_file_number, |
1565 | | const PreparedFileInfo* prepared_file_info, |
1566 | 0 | IngestedFileInfo* file_to_ingest, SuperVersion* sv) { |
1567 | 0 | file_to_ingest->external_file_path = external_file; |
1568 | |
|
1569 | 0 | std::unique_ptr<TableReader> table_reader; |
1570 | 0 | Status status = |
1571 | 0 | prepared_file_info != nullptr |
1572 | 0 | ? GetIngestedFileInfoFromFileInfo(external_file, *prepared_file_info, |
1573 | 0 | file_to_ingest) |
1574 | 0 | : GetIngestedFileInfoFromFile(external_file, new_file_number, |
1575 | 0 | file_to_ingest, sv, &table_reader); |
1576 | 0 | if (!status.ok()) { |
1577 | 0 | return status; |
1578 | 0 | } |
1579 | | |
1580 | 0 | assert(file_to_ingest->file_size > 0); |
1581 | 0 | assert(!file_to_ingest->unset()); |
1582 | 0 | assert(file_to_ingest->table_properties.num_entries > 0 || |
1583 | 0 | file_to_ingest->table_properties.num_range_deletions > 0); |
1584 | 0 | if (ingestion_options_.allow_db_generated_files) { |
1585 | | // These files keep their original sequence numbers (derived, not |
1586 | | // reassigned), so the bounds must be valid here. |
1587 | 0 | assert(file_to_ingest->smallest_seqno <= file_to_ingest->largest_seqno); |
1588 | 0 | assert(file_to_ingest->largest_seqno < kMaxSequenceNumber); |
1589 | 0 | } |
1590 | | |
1591 | | // Verify the file checksum if requested. The open-and-scan path already has a |
1592 | | // `TableReader` open; the fast-path opens one here only when verification is |
1593 | | // requested -- otherwise it performs no file I/O. |
1594 | 0 | if (ingestion_options_.verify_checksums_before_ingest) { |
1595 | 0 | if (table_reader == nullptr) { |
1596 | 0 | status = |
1597 | 0 | ResetTableReader(external_file, new_file_number, |
1598 | 0 | file_to_ingest->user_defined_timestamps_persisted, |
1599 | 0 | sv, file_to_ingest, &table_reader); |
1600 | 0 | if (!status.ok()) { |
1601 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
1602 | 0 | "Failed to reset table reader for external file %s: %s", |
1603 | 0 | external_file.c_str(), status.ToString().c_str()); |
1604 | 0 | return status; |
1605 | 0 | } |
1606 | 0 | } |
1607 | | // If customized readahead size is needed, we can pass a user option all the |
1608 | | // way to here. Right now we just rely on the default readahead to keep |
1609 | | // things simple. |
1610 | | // TODO: plumb Env::IOActivity, Env::IOPriority |
1611 | 0 | ReadOptions ro; |
1612 | 0 | ro.readahead_size = ingestion_options_.verify_checksums_readahead_size; |
1613 | 0 | ro.fill_cache = ingestion_options_.fill_cache; |
1614 | 0 | status = table_reader->VerifyChecksum( |
1615 | 0 | ro, TableReaderCaller::kExternalSSTIngestion); |
1616 | 0 | if (!status.ok()) { |
1617 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
1618 | 0 | "Failed to verify checksum for external file %s: %s", |
1619 | 0 | external_file.c_str(), status.ToString().c_str()); |
1620 | 0 | return status; |
1621 | 0 | } |
1622 | 0 | } |
1623 | | |
1624 | | // Assign FD with number. |
1625 | 0 | file_to_ingest->fd = |
1626 | 0 | FileDescriptor(new_file_number, 0, file_to_ingest->file_size); |
1627 | |
|
1628 | 0 | const size_t ts_sz = ucmp_->timestamp_size(); |
1629 | 0 | Slice smallest = file_to_ingest->smallest_internal_key.user_key(); |
1630 | 0 | Slice largest = file_to_ingest->largest_internal_key.user_key(); |
1631 | 0 | if (ts_sz > 0) { |
1632 | 0 | AppendUserKeyWithMaxTimestamp(&file_to_ingest->start_ukey, smallest, ts_sz); |
1633 | 0 | AppendUserKeyWithMinTimestamp(&file_to_ingest->limit_ukey, largest, ts_sz); |
1634 | 0 | } else { |
1635 | 0 | file_to_ingest->start_ukey.assign(smallest.data(), smallest.size()); |
1636 | 0 | file_to_ingest->limit_ukey.assign(largest.data(), largest.size()); |
1637 | 0 | } |
1638 | |
|
1639 | 0 | auto s = |
1640 | 0 | GetSstInternalUniqueId(file_to_ingest->table_properties.db_id, |
1641 | 0 | file_to_ingest->table_properties.db_session_id, |
1642 | 0 | file_to_ingest->table_properties.orig_file_number, |
1643 | 0 | &(file_to_ingest->unique_id)); |
1644 | 0 | if (!s.ok()) { |
1645 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
1646 | 0 | "Failed to get SST unique id for file %s", |
1647 | 0 | file_to_ingest->internal_file_path.c_str()); |
1648 | 0 | file_to_ingest->unique_id = kNullUniqueId64x2; |
1649 | 0 | } |
1650 | |
|
1651 | 0 | return status; |
1652 | 0 | } |
1653 | | |
1654 | | Status ExternalSstFileIngestionJob::AssignLevelAndSeqnoForIngestedFile( |
1655 | | SuperVersion* sv, bool force_global_seqno, CompactionStyle compaction_style, |
1656 | | SequenceNumber last_seqno, IngestedFileInfo* file_to_ingest, |
1657 | | SequenceNumber* assigned_seqno, |
1658 | 0 | std::optional<int> prev_batch_uppermost_level) { |
1659 | 0 | Status status; |
1660 | 0 | *assigned_seqno = 0; |
1661 | 0 | const size_t ts_sz = ucmp_->timestamp_size(); |
1662 | 0 | assert(!prev_batch_uppermost_level.has_value() || |
1663 | 0 | prev_batch_uppermost_level.value() < cfd_->NumberLevels()); |
1664 | 0 | bool must_assign_to_l0 = (prev_batch_uppermost_level.has_value() && |
1665 | 0 | prev_batch_uppermost_level.value() == 0) || |
1666 | 0 | compaction_style == kCompactionStyleFIFO; |
1667 | |
|
1668 | 0 | if (force_global_seqno || (!ingestion_options_.allow_db_generated_files && |
1669 | 0 | (files_overlap_ || must_assign_to_l0))) { |
1670 | 0 | *assigned_seqno = last_seqno + 1; |
1671 | 0 | if (must_assign_to_l0) { |
1672 | 0 | assert(ts_sz == 0); |
1673 | 0 | file_to_ingest->picked_level = 0; |
1674 | 0 | if (ingestion_options_.fail_if_not_bottommost_level && |
1675 | 0 | cfd_->NumberLevels() > 1) { |
1676 | 0 | status = Status::TryAgain( |
1677 | 0 | "Files cannot be ingested to Lmax. Please make sure key range of " |
1678 | 0 | "Lmax does not overlap with files to ingest."); |
1679 | 0 | } |
1680 | 0 | return status; |
1681 | 0 | } |
1682 | 0 | } |
1683 | | |
1684 | 0 | bool overlap_with_db = false; |
1685 | 0 | Arena arena; |
1686 | | // TODO: plumb Env::IOActivity, Env::IOPriority |
1687 | 0 | ReadOptions ro; |
1688 | 0 | ro.fill_cache = ingestion_options_.fill_cache; |
1689 | 0 | ro.total_order_seek = true; |
1690 | 0 | int target_level = 0; |
1691 | 0 | auto* vstorage = cfd_->current()->storage_info(); |
1692 | 0 | assert(!must_assign_to_l0 || ingestion_options_.allow_db_generated_files); |
1693 | 0 | int assigned_level_exclusive_end = cfd_->NumberLevels(); |
1694 | 0 | if (must_assign_to_l0) { |
1695 | 0 | assigned_level_exclusive_end = 0; |
1696 | 0 | } else if (prev_batch_uppermost_level.has_value()) { |
1697 | 0 | assigned_level_exclusive_end = prev_batch_uppermost_level.value(); |
1698 | 0 | } |
1699 | | |
1700 | | // When ingesting db generated files, we require that ingested files do not |
1701 | | // overlap with any file in the DB. So we need to check all levels. |
1702 | 0 | int overlap_checking_exclusive_end = |
1703 | 0 | ingestion_options_.allow_db_generated_files |
1704 | 0 | ? cfd_->NumberLevels() |
1705 | 0 | : assigned_level_exclusive_end; |
1706 | 0 | const bool is_generated_range_tombstone = |
1707 | 0 | atomic_replace_range_tombstone_active_ && |
1708 | 0 | file_to_ingest->generated_for_ingestion; |
1709 | 0 | for (int lvl = 0; lvl < overlap_checking_exclusive_end; lvl++) { |
1710 | 0 | if (lvl > 0 && lvl < vstorage->base_level()) { |
1711 | 0 | continue; |
1712 | 0 | } |
1713 | 0 | if (lvl < assigned_level_exclusive_end && |
1714 | 0 | atomic_replace_range_.has_value() && |
1715 | 0 | !atomic_replace_range_tombstone_active_) { |
1716 | 0 | target_level = lvl; |
1717 | 0 | continue; |
1718 | 0 | } |
1719 | 0 | if (cfd_->RangeOverlapWithCompaction(file_to_ingest->start_ukey, |
1720 | 0 | file_to_ingest->limit_ukey, lvl, |
1721 | | /*range_limit_exclusive=*/ |
1722 | 0 | is_generated_range_tombstone)) { |
1723 | | // We must use L0 or any level higher than `lvl` to be able to overwrite |
1724 | | // the compaction output keys that we overlap with in this level, We also |
1725 | | // need to assign this file a seqno to overwrite the compaction output |
1726 | | // keys in level `lvl` |
1727 | 0 | overlap_with_db = true; |
1728 | 0 | break; |
1729 | 0 | } else if (vstorage->NumLevelFiles(lvl) > 0) { |
1730 | 0 | bool overlap_with_level = false; |
1731 | 0 | status = sv->current->OverlapWithLevelIterator( |
1732 | 0 | ro, env_options_, file_to_ingest->start_ukey, |
1733 | 0 | file_to_ingest->limit_ukey, lvl, &overlap_with_level); |
1734 | 0 | if (!status.ok()) { |
1735 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
1736 | 0 | "Failed to check overlap with level iterator: %s", |
1737 | 0 | status.ToString().c_str()); |
1738 | 0 | return status; |
1739 | 0 | } |
1740 | 0 | if (overlap_with_level) { |
1741 | | // We must use L0 or any level higher than `lvl` to be able to overwrite |
1742 | | // the keys that we overlap with in this level, We also need to assign |
1743 | | // this file a seqno to overwrite the existing keys in level `lvl` |
1744 | 0 | overlap_with_db = true; |
1745 | 0 | break; |
1746 | 0 | } |
1747 | 0 | } |
1748 | | |
1749 | | // We don't overlap with any keys in this level, but we still need to check |
1750 | | // if our file can fit in it |
1751 | 0 | if (lvl < assigned_level_exclusive_end && |
1752 | 0 | IngestedFileFitInLevel(file_to_ingest, lvl)) { |
1753 | 0 | target_level = lvl; |
1754 | 0 | } |
1755 | 0 | } |
1756 | | |
1757 | 0 | if (ingestion_options_.fail_if_not_bottommost_level && |
1758 | 0 | target_level < cfd_->NumberLevels() - 1) { |
1759 | 0 | status = Status::TryAgain( |
1760 | 0 | "Files cannot be ingested to Lmax. Please make sure key range of Lmax " |
1761 | 0 | "and ongoing compaction's output to Lmax does not overlap with files " |
1762 | 0 | "to ingest. Input files overlapping with each other can cause some " |
1763 | 0 | "file to be assigned to non Lmax level."); |
1764 | 0 | return status; |
1765 | 0 | } |
1766 | | |
1767 | 0 | TEST_SYNC_POINT_CALLBACK( |
1768 | 0 | "ExternalSstFileIngestionJob::AssignLevelAndSeqnoForIngestedFile", |
1769 | 0 | &overlap_with_db); |
1770 | 0 | file_to_ingest->picked_level = target_level; |
1771 | 0 | if (overlap_with_db) { |
1772 | 0 | if (ts_sz > 0) { |
1773 | 0 | status = Status::InvalidArgument( |
1774 | 0 | "Column family enables user-defined timestamps, please make sure the " |
1775 | 0 | "key range (without timestamp) of external file does not overlap " |
1776 | 0 | "with key range (without timestamp) in the db"); |
1777 | 0 | return status; |
1778 | 0 | } |
1779 | 0 | if (*assigned_seqno == 0) { |
1780 | 0 | *assigned_seqno = last_seqno + 1; |
1781 | 0 | } |
1782 | 0 | } |
1783 | | |
1784 | 0 | return status; |
1785 | 0 | } |
1786 | | |
1787 | | Status ExternalSstFileIngestionJob::CheckLevelForIngestedBehindFile( |
1788 | 0 | IngestedFileInfo* file_to_ingest) { |
1789 | 0 | assert(!atomic_replace_range_.has_value()); |
1790 | |
|
1791 | 0 | auto* vstorage = cfd_->current()->storage_info(); |
1792 | | // First, check if new files fit in the last level |
1793 | 0 | int last_lvl = cfd_->NumberLevels() - 1; |
1794 | 0 | if (!IngestedFileFitInLevel(file_to_ingest, last_lvl)) { |
1795 | 0 | return Status::InvalidArgument( |
1796 | 0 | "Can't ingest_behind file as it doesn't fit " |
1797 | 0 | "at the last level!"); |
1798 | 0 | } |
1799 | | |
1800 | | // Second, check if despite cf_allow_ingest_behind=true we still have 0 |
1801 | | // seqnums at some upper level |
1802 | 0 | for (int lvl = 0; lvl < cfd_->NumberLevels() - 1; lvl++) { |
1803 | 0 | for (auto file : vstorage->LevelFiles(lvl)) { |
1804 | 0 | if (file->fd.smallest_seqno == 0) { |
1805 | 0 | return Status::InvalidArgument( |
1806 | 0 | "Can't ingest_behind file as despite cf_allow_ingest_behind=true " |
1807 | 0 | "there are files with 0 seqno in database at upper levels!"); |
1808 | 0 | } |
1809 | 0 | } |
1810 | 0 | } |
1811 | | |
1812 | 0 | file_to_ingest->picked_level = last_lvl; |
1813 | 0 | return Status::OK(); |
1814 | 0 | } |
1815 | | |
1816 | | Status ExternalSstFileIngestionJob::AssignGlobalSeqnoForIngestedFile( |
1817 | 0 | IngestedFileInfo* file_to_ingest, SequenceNumber seqno) { |
1818 | 0 | if (ingestion_options_.allow_db_generated_files) { |
1819 | 0 | assert(seqno == 0); |
1820 | 0 | assert(file_to_ingest->original_seqno == 0); |
1821 | 0 | } |
1822 | 0 | if (file_to_ingest->original_seqno == seqno) { |
1823 | | // This file already has the correct global seqno. |
1824 | 0 | return Status::OK(); |
1825 | 0 | } else if (!ingestion_options_.allow_global_seqno) { |
1826 | 0 | return Status::InvalidArgument("Global seqno is required, but disabled"); |
1827 | 0 | } else if (ingestion_options_.write_global_seqno && |
1828 | 0 | file_to_ingest->global_seqno_offset == 0) { |
1829 | 0 | return Status::InvalidArgument( |
1830 | 0 | "Trying to set global seqno for a file that don't have a global seqno " |
1831 | 0 | "field"); |
1832 | 0 | } |
1833 | | |
1834 | 0 | if (ingestion_options_.write_global_seqno) { |
1835 | | // Determine if we can write global_seqno to a given offset of file. |
1836 | | // If the file system does not support random write, then we should not. |
1837 | | // Otherwise we should. |
1838 | 0 | std::unique_ptr<FSRandomRWFile> rwfile; |
1839 | 0 | Status status = fs_->NewRandomRWFile(file_to_ingest->internal_file_path, |
1840 | 0 | env_options_, &rwfile, nullptr); |
1841 | 0 | TEST_SYNC_POINT_CALLBACK("ExternalSstFileIngestionJob::NewRandomRWFile", |
1842 | 0 | &status); |
1843 | 0 | if (status.ok()) { |
1844 | 0 | FSRandomRWFilePtr fsptr(std::move(rwfile), io_tracer_, |
1845 | 0 | file_to_ingest->internal_file_path); |
1846 | 0 | std::string seqno_val; |
1847 | 0 | PutFixed64(&seqno_val, seqno); |
1848 | 0 | status = fsptr->Write(file_to_ingest->global_seqno_offset, seqno_val, |
1849 | 0 | IOOptions(), nullptr); |
1850 | 0 | if (!status.ok()) { |
1851 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
1852 | 0 | "Failed to write global seqno to %s: %s", |
1853 | 0 | file_to_ingest->internal_file_path.c_str(), |
1854 | 0 | status.ToString().c_str()); |
1855 | 0 | return status; |
1856 | 0 | } |
1857 | | |
1858 | 0 | if (status.ok()) { |
1859 | 0 | TEST_SYNC_POINT("ExternalSstFileIngestionJob::BeforeSyncGlobalSeqno"); |
1860 | 0 | status = SyncIngestedFile(fsptr.get()); |
1861 | 0 | TEST_SYNC_POINT("ExternalSstFileIngestionJob::AfterSyncGlobalSeqno"); |
1862 | 0 | if (!status.ok()) { |
1863 | 0 | ROCKS_LOG_WARN(db_options_.info_log, |
1864 | 0 | "Failed to sync ingested file %s after writing global " |
1865 | 0 | "sequence number: %s", |
1866 | 0 | file_to_ingest->internal_file_path.c_str(), |
1867 | 0 | status.ToString().c_str()); |
1868 | 0 | } |
1869 | 0 | } |
1870 | 0 | if (!status.ok()) { |
1871 | 0 | return status; |
1872 | 0 | } |
1873 | 0 | } else if (!status.IsNotSupported()) { |
1874 | 0 | ROCKS_LOG_WARN( |
1875 | 0 | db_options_.info_log, |
1876 | 0 | "Failed to open ingested file %s for random read/write: %s", |
1877 | 0 | file_to_ingest->internal_file_path.c_str(), |
1878 | 0 | status.ToString().c_str()); |
1879 | 0 | return status; |
1880 | 0 | } |
1881 | 0 | } |
1882 | | |
1883 | 0 | file_to_ingest->assigned_seqno = seqno; |
1884 | 0 | return Status::OK(); |
1885 | 0 | } |
1886 | | |
1887 | | IOStatus ExternalSstFileIngestionJob::GenerateChecksumForIngestedFile( |
1888 | 0 | IngestedFileInfo* file_to_ingest) { |
1889 | 0 | if (db_options_.file_checksum_gen_factory == nullptr || |
1890 | 0 | need_generate_file_checksum_ == false || |
1891 | 0 | ingestion_options_.write_global_seqno == false) { |
1892 | | // If file_checksum_gen_factory is not set, we are not able to generate |
1893 | | // the checksum. if write_global_seqno is false, it means we will use |
1894 | | // file checksum generated during Prepare(). This step will be skipped. |
1895 | 0 | return IOStatus::OK(); |
1896 | 0 | } |
1897 | 0 | std::string file_checksum; |
1898 | 0 | std::string file_checksum_func_name; |
1899 | 0 | std::string requested_checksum_func_name; |
1900 | | // TODO: rate limit file reads for checksum calculation during file ingestion. |
1901 | | // TODO: plumb Env::IOActivity |
1902 | 0 | ReadOptions ro; |
1903 | 0 | FileOptions gen_fopts; |
1904 | 0 | gen_fopts.file_checksum_func_name = kNoFileChecksumFuncName; |
1905 | 0 | IOStatus io_s = GenerateOneFileChecksum( |
1906 | 0 | fs_.get(), file_to_ingest->internal_file_path, |
1907 | 0 | db_options_.file_checksum_gen_factory.get(), requested_checksum_func_name, |
1908 | 0 | &file_checksum, &file_checksum_func_name, |
1909 | 0 | ingestion_options_.verify_checksums_readahead_size, |
1910 | 0 | db_options_.allow_mmap_reads, io_tracer_, db_options_.rate_limiter.get(), |
1911 | 0 | ro, db_options_.stats, db_options_.clock, gen_fopts); |
1912 | 0 | if (!io_s.ok()) { |
1913 | 0 | ROCKS_LOG_WARN( |
1914 | 0 | db_options_.info_log, "Failed to generate checksum for %s: %s", |
1915 | 0 | file_to_ingest->internal_file_path.c_str(), io_s.ToString().c_str()); |
1916 | 0 | return io_s; |
1917 | 0 | } |
1918 | 0 | file_to_ingest->file_checksum = std::move(file_checksum); |
1919 | 0 | file_to_ingest->file_checksum_func_name = std::move(file_checksum_func_name); |
1920 | 0 | return IOStatus::OK(); |
1921 | 0 | } |
1922 | | |
1923 | | bool ExternalSstFileIngestionJob::IngestedFileFitInLevel( |
1924 | 0 | const IngestedFileInfo* file_to_ingest, int level) { |
1925 | 0 | if (level == 0) { |
1926 | | // Files can always fit in L0 |
1927 | 0 | return true; |
1928 | 0 | } |
1929 | | |
1930 | 0 | auto* vstorage = cfd_->current()->storage_info(); |
1931 | 0 | Slice file_smallest_user_key(file_to_ingest->start_ukey); |
1932 | 0 | Slice file_largest_user_key(file_to_ingest->limit_ukey); |
1933 | |
|
1934 | 0 | if (vstorage->OverlapInLevel(level, &file_smallest_user_key, |
1935 | 0 | &file_largest_user_key)) { |
1936 | | // File overlap with another files in this level, we cannot |
1937 | | // add it to this level |
1938 | 0 | return false; |
1939 | 0 | } |
1940 | | |
1941 | | // File did not overlap with level files, nor compaction output |
1942 | 0 | return true; |
1943 | 0 | } |
1944 | | |
1945 | | template <typename TWritableFile> |
1946 | 0 | Status ExternalSstFileIngestionJob::SyncIngestedFile(TWritableFile* file) { |
1947 | 0 | assert(file != nullptr); |
1948 | 0 | if (db_options_.use_fsync) { |
1949 | 0 | return file->Fsync(IOOptions(), nullptr); |
1950 | 0 | } else { |
1951 | 0 | return file->Sync(IOOptions(), nullptr); |
1952 | 0 | } |
1953 | 0 | } |
1954 | | |
1955 | | Status ExternalSstFileIngestionJob::GetSeqnoBoundaryForFile( |
1956 | | TableReader* table_reader, SuperVersion* sv, |
1957 | 0 | IngestedFileInfo* file_to_ingest, bool allow_data_in_errors) { |
1958 | 0 | const auto tp = table_reader->GetTableProperties(); |
1959 | 0 | const bool has_largest_seqno = tp->HasKeyLargestSeqno(); |
1960 | 0 | SequenceNumber largest_seqno = tp->key_largest_seqno; |
1961 | 0 | if (has_largest_seqno) { |
1962 | 0 | file_to_ingest->largest_seqno = largest_seqno; |
1963 | 0 | if (largest_seqno == 0) { |
1964 | 0 | file_to_ingest->smallest_seqno = 0; |
1965 | 0 | return Status::OK(); |
1966 | 0 | } |
1967 | 0 | if (tp->HasKeySmallestSeqno()) { |
1968 | 0 | file_to_ingest->smallest_seqno = tp->key_smallest_seqno; |
1969 | 0 | return Status::OK(); |
1970 | 0 | } |
1971 | 0 | } |
1972 | | |
1973 | | // For older SST files they may not be recorded in table properties, so |
1974 | | // we scan the file to find out. |
1975 | 0 | TEST_SYNC_POINT( |
1976 | 0 | "ExternalSstFileIngestionJob::GetSeqnoBoundaryForFile:FileScan"); |
1977 | 0 | SequenceNumber smallest_seqno = kMaxSequenceNumber; |
1978 | 0 | SequenceNumber largest_seqno_from_iter = 0; |
1979 | 0 | ReadOptions ro; |
1980 | 0 | ro.fill_cache = ingestion_options_.fill_cache; |
1981 | 0 | std::unique_ptr<InternalIterator> iter(table_reader->NewIterator( |
1982 | 0 | ro, sv->mutable_cf_options.prefix_extractor.get(), /*arena=*/nullptr, |
1983 | 0 | /*skip_filters=*/false, TableReaderCaller::kExternalSSTIngestion)); |
1984 | 0 | ParsedInternalKey key; |
1985 | 0 | iter->SeekToFirst(); |
1986 | 0 | while (iter->Valid()) { |
1987 | 0 | Status pik_status = |
1988 | 0 | ParseInternalKey(iter->key(), &key, allow_data_in_errors); |
1989 | 0 | if (!pik_status.ok()) { |
1990 | 0 | return Status::Corruption("Corrupted key in external file. ", |
1991 | 0 | pik_status.getState()); |
1992 | 0 | } |
1993 | 0 | smallest_seqno = std::min(smallest_seqno, key.sequence); |
1994 | 0 | largest_seqno_from_iter = std::max(largest_seqno_from_iter, key.sequence); |
1995 | 0 | iter->Next(); |
1996 | 0 | } |
1997 | 0 | if (!iter->status().ok()) { |
1998 | 0 | return iter->status(); |
1999 | 0 | } |
2000 | | |
2001 | 0 | if (table_reader->GetTableProperties()->num_range_deletions > 0) { |
2002 | 0 | std::unique_ptr<InternalIterator> range_del_iter( |
2003 | 0 | table_reader->NewRangeTombstoneIterator(ro)); |
2004 | 0 | if (range_del_iter != nullptr) { |
2005 | 0 | for (range_del_iter->SeekToFirst(); range_del_iter->Valid(); |
2006 | 0 | range_del_iter->Next()) { |
2007 | 0 | Status pik_status = |
2008 | 0 | ParseInternalKey(range_del_iter->key(), &key, allow_data_in_errors); |
2009 | 0 | if (!pik_status.ok()) { |
2010 | 0 | return Status::Corruption("Corrupted key in external file. ", |
2011 | 0 | pik_status.getState()); |
2012 | 0 | } |
2013 | 0 | smallest_seqno = std::min(smallest_seqno, key.sequence); |
2014 | 0 | largest_seqno_from_iter = |
2015 | 0 | std::max(largest_seqno_from_iter, key.sequence); |
2016 | 0 | } |
2017 | 0 | if (!range_del_iter->status().ok()) { |
2018 | 0 | return range_del_iter->status(); |
2019 | 0 | } |
2020 | 0 | } |
2021 | 0 | } |
2022 | | |
2023 | 0 | file_to_ingest->smallest_seqno = smallest_seqno; |
2024 | 0 | if (!has_largest_seqno) { |
2025 | 0 | file_to_ingest->largest_seqno = largest_seqno_from_iter; |
2026 | 0 | } else { |
2027 | 0 | assert(largest_seqno == largest_seqno_from_iter); |
2028 | 0 | file_to_ingest->largest_seqno = largest_seqno; |
2029 | 0 | } |
2030 | |
|
2031 | 0 | if (file_to_ingest->largest_seqno == kMaxSequenceNumber) { |
2032 | 0 | return Status::InvalidArgument( |
2033 | 0 | "Unknown smallest seqno for db generated file."); |
2034 | 0 | } |
2035 | 0 | if (file_to_ingest->smallest_seqno == kMaxSequenceNumber) { |
2036 | 0 | return Status::InvalidArgument( |
2037 | 0 | "Unknown largest seqno for db generated file."); |
2038 | 0 | } |
2039 | 0 | return Status::OK(); |
2040 | 0 | } |
2041 | | |
2042 | | } // namespace ROCKSDB_NAMESPACE |