Skip to main content

mito2/region/
utils.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::sync::Arc;
16use std::time::Instant;
17
18use common_base::readable_size::ReadableSize;
19use common_telemetry::{debug, error, info};
20use futures::future::try_join_all;
21use object_store::manager::ObjectStoreManagerRef;
22use snafu::{ResultExt, ensure};
23use store_api::metadata::RegionMetadataRef;
24use store_api::region_request::PathType;
25use store_api::storage::{FileId, IndexVersion, RegionId};
26
27use crate::access_layer::AccessLayerRef;
28use crate::config::MitoConfig;
29use crate::error::{self, InvalidSourceAndTargetRegionSnafu, Result};
30use crate::manifest::action::RegionManifest;
31use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions};
32use crate::region::opener::get_object_store;
33use crate::region::options::RegionOptions;
34use crate::sst::file::{RegionFileId, RegionIndexId};
35use crate::sst::location;
36
37/// A loader for loading metadata from a region dir.
38#[derive(Debug, Clone)]
39pub struct RegionMetadataLoader {
40    config: Arc<MitoConfig>,
41    object_store_manager: ObjectStoreManagerRef,
42}
43
44impl RegionMetadataLoader {
45    /// Creates a new `RegionMetadataLoader`.
46    pub fn new(config: Arc<MitoConfig>, object_store_manager: ObjectStoreManagerRef) -> Self {
47        Self {
48            config,
49            object_store_manager,
50        }
51    }
52
53    /// Loads the metadata of the region from the region dir.
54    pub async fn load(
55        &self,
56        region_dir: &str,
57        region_options: &RegionOptions,
58    ) -> Result<Option<RegionMetadataRef>> {
59        let manifest = self
60            .load_manifest(region_dir, &region_options.storage)
61            .await?;
62        Ok(manifest.map(|m| m.metadata.clone()))
63    }
64
65    /// Loads the manifest of the region from the region dir.
66    pub async fn load_manifest(
67        &self,
68        region_dir: &str,
69        storage: &Option<String>,
70    ) -> Result<Option<Arc<RegionManifest>>> {
71        let object_store = get_object_store(storage, &self.object_store_manager)?;
72        let region_manifest_options =
73            RegionManifestOptions::new(&self.config, region_dir, &object_store);
74        let Some(manifest_manager) =
75            RegionManifestManager::open(region_manifest_options, &Default::default()).await?
76        else {
77            return Ok(None);
78        };
79
80        let manifest = manifest_manager.manifest();
81        Ok(Some(manifest))
82    }
83}
84
85/// A copier for copying files from a region to another region.
86#[derive(Debug, Clone)]
87pub struct RegionFileCopier {
88    access_layer: AccessLayerRef,
89}
90
91/// A descriptor for a file.
92#[derive(Debug, Clone, Copy)]
93pub enum FileDescriptor {
94    /// An index file.
95    Index {
96        file_id: FileId,
97        version: IndexVersion,
98        size: u64,
99    },
100    /// A data file.
101    Data { file_id: FileId, size: u64 },
102}
103
104impl FileDescriptor {
105    pub fn size(&self) -> u64 {
106        match self {
107            FileDescriptor::Index { size, .. } => *size,
108            FileDescriptor::Data { size, .. } => *size,
109        }
110    }
111
112    fn file_path(&self, region_id: RegionId, table_dir: &str, path_type: PathType) -> String {
113        match *self {
114            FileDescriptor::Index {
115                file_id, version, ..
116            } => location::index_file_path(
117                table_dir,
118                RegionIndexId::new(RegionFileId::new(region_id, file_id), version),
119                path_type,
120            ),
121            FileDescriptor::Data { file_id, .. } => {
122                location::sst_file_path(table_dir, RegionFileId::new(region_id, file_id), path_type)
123            }
124        }
125    }
126}
127
128/// Builds the source and target file paths for a given file descriptor.
129///
130/// # Arguments
131///
132/// * `source_region_id`: The ID of the source region.
133/// * `target_region_id`: The ID of the target region.
134/// * `file_id`: The ID of the file.
135///
136/// # Returns
137///
138/// A tuple containing the source and target file paths.
139fn build_copy_file_paths(
140    source_region_id: RegionId,
141    target_region_id: RegionId,
142    file_descriptor: FileDescriptor,
143    table_dir: &str,
144    path_type: PathType,
145) -> (String, String) {
146    (
147        file_descriptor.file_path(source_region_id, table_dir, path_type),
148        file_descriptor.file_path(target_region_id, table_dir, path_type),
149    )
150}
151
152fn build_delete_file_path(
153    target_region_id: RegionId,
154    file_descriptor: FileDescriptor,
155    table_dir: &str,
156    path_type: PathType,
157) -> String {
158    file_descriptor.file_path(target_region_id, table_dir, path_type)
159}
160
161impl RegionFileCopier {
162    pub fn new(access_layer: AccessLayerRef) -> Self {
163        Self { access_layer }
164    }
165
166    /// Copies files from a source region to a target region.
167    ///
168    /// # Arguments
169    ///
170    /// * `source_region_id`: The ID of the source region.
171    /// * `target_region_id`: The ID of the target region.
172    /// * `file_ids`: The IDs of the files to copy.
173    pub async fn copy_files(
174        &self,
175        source_region_id: RegionId,
176        target_region_id: RegionId,
177        file_ids: Vec<FileDescriptor>,
178        parallelism: usize,
179    ) -> Result<()> {
180        ensure!(
181            source_region_id.table_id() == target_region_id.table_id(),
182            InvalidSourceAndTargetRegionSnafu {
183                source_region_id,
184                target_region_id,
185            },
186        );
187        let table_dir = self.access_layer.table_dir();
188        let path_type = self.access_layer.path_type();
189        let object_store = self.access_layer.object_store();
190
191        info!(
192            "Copying {} files from region {} to region {}",
193            file_ids.len(),
194            source_region_id,
195            target_region_id
196        );
197        debug!(
198            "Copying files: {:?} from region {} to region {}",
199            file_ids, source_region_id, target_region_id
200        );
201        let mut tasks = Vec::with_capacity(parallelism);
202        for skip in 0..parallelism {
203            let target_file_ids = file_ids.iter().skip(skip).step_by(parallelism).copied();
204            let object_store = object_store.clone();
205            tasks.push(async move {
206                for file_desc in target_file_ids {
207                    let (source_path, target_path) = build_copy_file_paths(
208                        source_region_id,
209                        target_region_id,
210                        file_desc,
211                        table_dir,
212                        path_type,
213                    );
214                    let now = Instant::now();
215                    object_store
216                        .copy(&source_path, &target_path)
217                        .await
218                        .inspect_err(
219                            |e| error!(e; "Failed to copy file {} to {}", source_path, target_path),
220                        )
221                        .context(error::OpenDalSnafu)?;
222                    let file_size = ReadableSize(file_desc.size());
223                    info!(
224                        "Copied file {} to {}, file size: {}, elapsed: {:?}",
225                        source_path,
226                        target_path,
227                        file_size,
228                        now.elapsed(),
229                    );
230                }
231
232                Ok(())
233            });
234        }
235
236        if let Err(err) = try_join_all(tasks).await {
237            error!(err; "Failed to copy files from region {} to region {}", source_region_id, target_region_id);
238            self.clean_target_region(target_region_id, file_ids).await;
239            return Err(err);
240        }
241
242        Ok(())
243    }
244
245    /// Cleans the copied files from the target region.
246    async fn clean_target_region(&self, target_region_id: RegionId, file_ids: Vec<FileDescriptor>) {
247        let table_dir = self.access_layer.table_dir();
248        let path_type = self.access_layer.path_type();
249        let object_store = self.access_layer.object_store();
250        let delete_file_path = file_ids
251            .into_iter()
252            .map(|file_descriptor| {
253                build_delete_file_path(target_region_id, file_descriptor, table_dir, path_type)
254            })
255            .collect::<Vec<_>>();
256        debug!(
257            "Deleting files: {:?} after failed to copy files to target region {}",
258            delete_file_path, target_region_id
259        );
260        if let Err(err) = object_store.delete_iter(delete_file_path).await {
261            error!(err; "Failed to delete files from region {}", target_region_id);
262        }
263    }
264}
265
266#[cfg(test)]
267mod tests {
268    #[cfg(feature = "hdfs-object-store")]
269    use object_store::ObjectStore;
270    #[cfg(feature = "hdfs-object-store")]
271    use object_store::layers::HdfsCompatibilityLayer;
272    #[cfg(feature = "hdfs-object-store")]
273    use object_store::services::Fs;
274
275    use super::*;
276    #[cfg(feature = "hdfs-object-store")]
277    use crate::access_layer::AccessLayer;
278    #[cfg(feature = "hdfs-object-store")]
279    use crate::sst::index::intermediate::IntermediateManager;
280    #[cfg(feature = "hdfs-object-store")]
281    use crate::sst::index::puffin_manager::PuffinManagerFactory;
282
283    #[test]
284    fn test_build_copy_file_paths() {
285        common_telemetry::init_default_ut_logging();
286        let file_id = FileId::random();
287        let source_region_id = RegionId::new(1, 1);
288        let target_region_id = RegionId::new(1, 2);
289        let file_descriptor = FileDescriptor::Data { file_id, size: 100 };
290        let table_dir = "/table_dir";
291        let path_type = PathType::Bare;
292        let (source_path, target_path) = build_copy_file_paths(
293            source_region_id,
294            target_region_id,
295            file_descriptor,
296            table_dir,
297            path_type,
298        );
299        assert_eq!(
300            source_path,
301            format!("/table_dir/1_0000000001/{}.parquet", file_id)
302        );
303        assert_eq!(
304            target_path,
305            format!("/table_dir/1_0000000002/{}.parquet", file_id)
306        );
307
308        let version = 1;
309        let file_descriptor = FileDescriptor::Index {
310            file_id,
311            version,
312            size: 100,
313        };
314        let (source_path, target_path) = build_copy_file_paths(
315            source_region_id,
316            target_region_id,
317            file_descriptor,
318            table_dir,
319            path_type,
320        );
321        assert_eq!(
322            source_path,
323            format!(
324                "/table_dir/1_0000000001/index/{}.{}.puffin",
325                file_id, version
326            )
327        );
328        assert_eq!(
329            target_path,
330            format!(
331                "/table_dir/1_0000000002/index/{}.{}.puffin",
332                file_id, version
333            )
334        );
335    }
336
337    #[test]
338    fn test_build_delete_file_path() {
339        common_telemetry::init_default_ut_logging();
340        let file_id = FileId::random();
341        let target_region_id = RegionId::new(1, 2);
342        let table_dir = "/table_dir";
343        let path_type = PathType::Bare;
344
345        let file_descriptor = FileDescriptor::Data { file_id, size: 100 };
346        let path = build_delete_file_path(target_region_id, file_descriptor, table_dir, path_type);
347        assert_eq!(path, format!("/table_dir/1_0000000002/{}.parquet", file_id));
348
349        let file_descriptor = FileDescriptor::Index {
350            file_id,
351            version: 1,
352            size: 100,
353        };
354        let path = build_delete_file_path(target_region_id, file_descriptor, table_dir, path_type);
355        assert_eq!(
356            path,
357            format!("/table_dir/1_0000000002/index/{}.1.puffin", file_id)
358        );
359    }
360
361    #[cfg(feature = "hdfs-object-store")]
362    #[tokio::test]
363    async fn test_copy_region_files_with_hdfs_fallback() {
364        let (temp_dir, puffin_manager) =
365            PuffinManagerFactory::new_for_test_async("hdfs-copy-region").await;
366        let intermediate_manager = IntermediateManager::init_fs(temp_dir.path().to_string_lossy())
367            .await
368            .unwrap();
369        let storage_dir = temp_dir.path().join("storage");
370        std::fs::create_dir(&storage_dir).unwrap();
371        let object_store = ObjectStore::new(Fs::default().root(storage_dir.to_str().unwrap()))
372            .unwrap()
373            .layer(HdfsCompatibilityLayer::new_for_test());
374        let access_layer = Arc::new(AccessLayer::new(
375            "table_dir",
376            PathType::Bare,
377            object_store.clone(),
378            puffin_manager,
379            intermediate_manager,
380        ));
381        let copier = RegionFileCopier::new(access_layer);
382        let source_region_id = RegionId::new(1, 1);
383        let target_region_id = RegionId::new(1, 2);
384        let file_id = FileId::random();
385        let descriptor = FileDescriptor::Data { file_id, size: 8 };
386        let (source_path, target_path) = build_copy_file_paths(
387            source_region_id,
388            target_region_id,
389            descriptor,
390            "table_dir",
391            PathType::Bare,
392        );
393        object_store.write(&source_path, "contents").await.unwrap();
394
395        copier
396            .copy_files(source_region_id, target_region_id, vec![descriptor], 1)
397            .await
398            .unwrap();
399
400        assert_eq!(
401            b"contents",
402            object_store
403                .read(&target_path)
404                .await
405                .unwrap()
406                .to_bytes()
407                .as_ref()
408        );
409    }
410}