Skip to main content

cmd/datanode/
objbench.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::path::PathBuf;
16use std::sync::Arc;
17use std::time::Instant;
18
19use clap::Parser;
20use colored::Colorize;
21use futures::stream;
22use mito2::access_layer::{
23    AccessLayer, AccessLayerRef, Metrics, OperationType, SstWriteRequest, WriteType,
24};
25use mito2::cache::{CacheManager, CacheManagerRef};
26use mito2::config::{FulltextIndexConfig, MitoConfig, Mode};
27use mito2::read::FlatSource;
28use mito2::sst::FormatType;
29use mito2::sst::file::{FileHandle, FileMeta};
30use mito2::sst::file_purger::{FilePurger, FilePurgerRef};
31use mito2::sst::index::intermediate::IntermediateManager;
32use mito2::sst::index::puffin_manager::PuffinManagerFactory;
33use mito2::sst::parquet::WriteOptions;
34use mito2::sst::parquet::reader::ParquetReaderBuilder;
35use mito2::worker::write_cache_from_config;
36use object_store::ObjectStore;
37use parquet::file::metadata::FooterTail;
38use regex::Regex;
39use snafu::OptionExt;
40use store_api::path_utils::region_name;
41use store_api::region_request::PathType;
42use store_api::storage::FileId;
43
44use crate::datanode::tool_util::{
45    build_object_store, extract_region_metadata, max_row_group_uncompressed_size, parse_config,
46};
47use crate::error;
48
49/// Object storage benchmark command
50#[derive(Debug, Parser)]
51pub struct ObjbenchCommand {
52    /// Path to the object-store config file (TOML). Must deserialize into object_store::config::ObjectStoreConfig.
53    #[clap(long, value_name = "FILE")]
54    pub config: PathBuf,
55
56    /// Source SST file path in object-store (e.g. "region_dir/<uuid>.parquet").
57    #[clap(long, value_name = "PATH")]
58    pub source: String,
59
60    /// Verbose output
61    #[clap(short, long, default_value_t = false)]
62    pub verbose: bool,
63
64    /// Output file path for pprof flamegraph (enables profiling)
65    #[clap(long, value_name = "FILE")]
66    pub pprof_file: Option<PathBuf>,
67}
68
69impl ObjbenchCommand {
70    pub async fn run(&self) -> error::Result<()> {
71        if self.verbose {
72            common_telemetry::init_default_ut_logging();
73        }
74
75        println!("{}", "Starting objbench with config:".cyan().bold());
76
77        // Build object store from config
78        let (store_cfg, mut mito_engine_config, _wal_config) = parse_config(&self.config)?;
79
80        let object_store = build_object_store(&store_cfg).await?;
81        println!("{} Object store initialized", "✓".green());
82
83        // Prepare source identifiers
84        let components = parse_file_dir_components(&self.source)?;
85        println!(
86            "{} Source path parsed: {}, components: {:?}",
87            "✓".green(),
88            self.source,
89            components
90        );
91
92        // Load parquet metadata to extract RegionMetadata and file stats
93        println!("{}", "Loading parquet metadata...".yellow());
94        let file_size = object_store
95            .stat(&self.source)
96            .await
97            .map_err(|e| {
98                error::IllegalConfigSnafu {
99                    msg: format!("stat failed: {e}"),
100                }
101                .build()
102            })?
103            .content_length();
104        let parquet_meta = load_parquet_metadata(object_store.clone(), &self.source, file_size)
105            .await
106            .map_err(|e| {
107                error::IllegalConfigSnafu {
108                    msg: format!("read parquet metadata failed: {e}"),
109                }
110                .build()
111            })?;
112
113        let region_meta = extract_region_metadata(&self.source, &parquet_meta)?;
114        let num_rows = parquet_meta.file_metadata().num_rows() as u64;
115        let num_row_groups = parquet_meta.num_row_groups() as u64;
116        let max_row_group_uncompressed_size = max_row_group_uncompressed_size(&parquet_meta);
117
118        println!(
119            "{} Metadata loaded - rows: {}, size: {} bytes",
120            "✓".green(),
121            num_rows,
122            file_size
123        );
124
125        // Build a FileHandle for the source file
126        let file_meta = FileMeta {
127            region_id: region_meta.region_id,
128            file_id: components.file_id,
129            time_range: Default::default(),
130            level: 0,
131            file_size,
132            max_row_group_uncompressed_size,
133            available_indexes: Default::default(),
134            indexes: Default::default(),
135            index_file_size: 0,
136            index_version: 0,
137            num_rows,
138            num_row_groups,
139            sequence: None,
140            partition_expr: None,
141            num_series: 0,
142            ..Default::default()
143        };
144        let src_handle = FileHandle::new(file_meta, new_noop_file_purger());
145
146        // Build the reader for a single file via ParquetReaderBuilder
147        let table_dir = components.table_dir();
148        let (src_access_layer, cache_manager) = build_access_layer_simple(
149            &components,
150            object_store.clone(),
151            &mut mito_engine_config,
152            &store_cfg.data_home,
153        )
154        .await?;
155        let reader_build_start = Instant::now();
156
157        let reader = ParquetReaderBuilder::new(
158            table_dir,
159            components.path_type,
160            src_handle.clone(),
161            object_store.clone(),
162        )
163        .expected_metadata(Some(region_meta.clone()))
164        .build()
165        .await
166        .map_err(|e| {
167            error::IllegalConfigSnafu {
168                msg: format!("build reader failed: {e:?}"),
169            }
170            .build()
171        })?;
172        let reader = reader.ok_or_else(|| {
173            error::IllegalConfigSnafu {
174                msg: format!(
175                    "build reader returned no readable rows for source file {}",
176                    src_handle.file_id()
177                ),
178            }
179            .build()
180        })?;
181
182        let reader_build_elapsed = reader_build_start.elapsed();
183        let total_rows = reader.parquet_metadata().file_metadata().num_rows();
184        println!("{} Reader built in {:?}", "✓".green(), reader_build_elapsed);
185        let reader_stream = Box::pin(stream::try_unfold(reader, |mut reader| async move {
186            let batch = reader.next_record_batch().await?;
187            Ok(batch.map(|batch| (batch, reader)))
188        }));
189
190        // Build write request
191        let fulltext_index_config = FulltextIndexConfig {
192            create_on_compaction: Mode::Disable,
193            ..Default::default()
194        };
195
196        let source =
197            FlatSource::new_stream(region_meta.schema.arrow_schema().clone(), reader_stream);
198        let write_req = SstWriteRequest {
199            op_type: OperationType::Flush,
200            metadata: region_meta,
201            source,
202            cache_manager,
203            storage: None,
204            max_sequence: None,
205            sst_write_format: FormatType::PrimaryKey,
206            index_options: Default::default(),
207            index_config: mito_engine_config.index.clone(),
208            inverted_index_config: MitoConfig::default().inverted_index,
209            fulltext_index_config,
210            bloom_filter_index_config: MitoConfig::default().bloom_filter_index,
211            preserve_row_sequence: false,
212            #[cfg(feature = "vector_index")]
213            vector_index_config: Default::default(),
214        };
215
216        // Write SST
217        println!("{}", "Writing SST...".yellow());
218
219        // Start profiling if pprof_file is specified
220        #[cfg(unix)]
221        let profiler_guard = if self.pprof_file.is_some() {
222            println!("{} Starting profiling...", "⚡".yellow());
223            Some(
224                pprof::ProfilerGuardBuilder::default()
225                    .frequency(99)
226                    .blocklist(&["libc", "libgcc", "pthread", "vdso"])
227                    .build()
228                    .map_err(|e| {
229                        error::IllegalConfigSnafu {
230                            msg: format!("Failed to start profiler: {e}"),
231                        }
232                        .build()
233                    })?,
234            )
235        } else {
236            None
237        };
238
239        #[cfg(not(unix))]
240        if self.pprof_file.is_some() {
241            eprintln!(
242                "{}: Profiling is not supported on this platform",
243                "Warning".yellow()
244            );
245        }
246
247        let write_start = Instant::now();
248        let mut metrics = Metrics::new(WriteType::Flush);
249        let infos = src_access_layer
250            .write_sst(write_req, &WriteOptions::default(), &mut metrics)
251            .await
252            .map_err(|e| {
253                error::IllegalConfigSnafu {
254                    msg: format!("write_sst failed: {e:?}"),
255                }
256                .build()
257            })?;
258
259        let write_elapsed = write_start.elapsed();
260
261        // Stop profiling and generate flamegraph if enabled
262        #[cfg(unix)]
263        if let (Some(guard), Some(pprof_file)) = (profiler_guard, &self.pprof_file) {
264            println!("{} Generating flamegraph...", "🔥".yellow());
265            match guard.report().build() {
266                Ok(report) => {
267                    let mut flamegraph_data = Vec::new();
268                    if let Err(e) = report.flamegraph(&mut flamegraph_data) {
269                        println!("{}: Failed to generate flamegraph: {}", "Error".red(), e);
270                    } else if let Err(e) = std::fs::write(pprof_file, flamegraph_data) {
271                        println!(
272                            "{}: Failed to write flamegraph to {}: {}",
273                            "Error".red(),
274                            pprof_file.display(),
275                            e
276                        );
277                    } else {
278                        println!(
279                            "{} Flamegraph saved to {}",
280                            "✓".green(),
281                            pprof_file.display().to_string().cyan()
282                        );
283                    }
284                }
285                Err(e) => {
286                    println!("{}: Failed to generate pprof report: {}", "Error".red(), e);
287                }
288            }
289        }
290        assert_eq!(infos.len(), 1);
291        let dst_file_id = infos[0].file_id;
292        let dst_file_path = format!("{}/{}.parquet", components.region_dir(), dst_file_id);
293        let mut dst_index_path = None;
294        if infos[0].index_metadata.file_size > 0 {
295            dst_index_path = Some(format!(
296                "{}/index/{}.puffin",
297                components.region_dir(),
298                dst_file_id
299            ));
300        }
301
302        // Report results with ANSI colors
303        println!("\n{} {}", "Write complete!".green().bold(), "✓".green());
304        println!("  {}: {}", "Destination file".bold(), dst_file_path.cyan());
305        println!("  {}: {}", "Rows".bold(), total_rows.to_string().cyan());
306        println!(
307            "  {}: {}",
308            "File size".bold(),
309            format!("{} bytes", file_size).cyan()
310        );
311        println!(
312            "  {}: {:?}",
313            "Reader build time".bold(),
314            reader_build_elapsed
315        );
316        println!("  {}: {:?}", "Total time".bold(), write_elapsed);
317
318        // Print metrics in a formatted way
319        println!("  {}: {:?}", "Metrics".bold(), metrics,);
320
321        // Print infos
322        println!("  {}: {:?}", "Index".bold(), infos[0].index_metadata);
323
324        // Cleanup
325        println!("\n{}", "Cleaning up...".yellow());
326        object_store.delete(&dst_file_path).await.map_err(|e| {
327            error::IllegalConfigSnafu {
328                msg: format!("Failed to delete dest file {}: {}", dst_file_path, e),
329            }
330            .build()
331        })?;
332        println!("{} Temporary file {} deleted", "✓".green(), dst_file_path);
333
334        if let Some(index_path) = dst_index_path {
335            object_store.delete(&index_path).await.map_err(|e| {
336                error::IllegalConfigSnafu {
337                    msg: format!("Failed to delete dest index file {}: {}", index_path, e),
338                }
339                .build()
340            })?;
341            println!(
342                "{} Temporary index file {} deleted",
343                "✓".green(),
344                index_path
345            );
346        }
347
348        println!("\n{}", "Benchmark completed successfully!".green().bold());
349        Ok(())
350    }
351}
352
353#[derive(Debug)]
354struct FileDirComponents {
355    catalog: String,
356    schema: String,
357    table_id: u32,
358    region_sequence: u32,
359    path_type: PathType,
360    file_id: FileId,
361}
362
363impl FileDirComponents {
364    fn table_dir(&self) -> String {
365        format!("data/{}/{}/{}", self.catalog, self.schema, self.table_id)
366    }
367
368    fn region_dir(&self) -> String {
369        let region_name = region_name(self.table_id, self.region_sequence);
370        match self.path_type {
371            PathType::Bare => {
372                format!(
373                    "data/{}/{}/{}/{}",
374                    self.catalog, self.schema, self.table_id, region_name
375                )
376            }
377            PathType::Data => {
378                format!(
379                    "data/{}/{}/{}/{}/data",
380                    self.catalog, self.schema, self.table_id, region_name
381                )
382            }
383            PathType::Metadata => {
384                format!(
385                    "data/{}/{}/{}/{}/metadata",
386                    self.catalog, self.schema, self.table_id, region_name
387                )
388            }
389        }
390    }
391}
392
393fn parse_file_dir_components(path: &str) -> error::Result<FileDirComponents> {
394    // Define the regex pattern to match all three path styles
395    let pattern =
396        r"^data/([^/]+)/([^/]+)/([^/]+)/([^/]+)_([^/]+)(?:/data|/metadata)?/(.+).parquet$";
397
398    // Compile the regex
399    let re = Regex::new(pattern).expect("Invalid regex pattern");
400
401    // Determine the path type
402    let path_type = if path.contains("/data/") {
403        PathType::Data
404    } else if path.contains("/metadata/") {
405        PathType::Metadata
406    } else {
407        PathType::Bare
408    };
409
410    // Try to match the path
411    let components = (|| {
412        let captures = re.captures(path)?;
413        if captures.len() != 7 {
414            return None;
415        }
416        let mut components = FileDirComponents {
417            catalog: "".to_string(),
418            schema: "".to_string(),
419            table_id: 0,
420            region_sequence: 0,
421            path_type,
422            file_id: FileId::default(),
423        };
424        // Extract the components
425        components.catalog = captures.get(1)?.as_str().to_string();
426        components.schema = captures.get(2)?.as_str().to_string();
427        components.table_id = captures[3].parse().ok()?;
428        components.region_sequence = captures[5].parse().ok()?;
429        let file_id_str = &captures[6];
430        components.file_id = FileId::parse_str(file_id_str).ok()?;
431        Some(components)
432    })();
433    components.context(error::IllegalConfigSnafu {
434        msg: format!("Expect valid source file path, got: {}", path),
435    })
436}
437
438async fn build_access_layer_simple(
439    components: &FileDirComponents,
440    object_store: ObjectStore,
441    config: &mut MitoConfig,
442    data_home: &str,
443) -> error::Result<(AccessLayerRef, CacheManagerRef)> {
444    let _ = config.index.sanitize(data_home, &config.inverted_index);
445    let puffin_manager = PuffinManagerFactory::new(
446        &config.index.aux_path,
447        config.index.staging_size.as_bytes(),
448        Some(config.index.write_buffer_size.as_bytes() as _),
449        config.index.staging_ttl,
450    )
451    .await
452    .map_err(|e| {
453        error::IllegalConfigSnafu {
454            msg: format!("Failed to build access layer: {e:?}"),
455        }
456        .build()
457    })?;
458
459    let intermediate_manager = IntermediateManager::init_fs(&config.index.aux_path)
460        .await
461        .map_err(|e| {
462            error::IllegalConfigSnafu {
463                msg: format!("Failed to build IntermediateManager: {e:?}"),
464            }
465            .build()
466        })?
467        .with_buffer_size(Some(config.index.write_buffer_size.as_bytes() as _));
468
469    let cache_manager =
470        build_cache_manager(config, puffin_manager.clone(), intermediate_manager.clone()).await?;
471    let layer = AccessLayer::new(
472        components.table_dir(),
473        components.path_type,
474        object_store,
475        puffin_manager,
476        intermediate_manager,
477    );
478    Ok((Arc::new(layer), cache_manager))
479}
480
481async fn build_cache_manager(
482    config: &MitoConfig,
483    puffin_manager: PuffinManagerFactory,
484    intermediate_manager: IntermediateManager,
485) -> error::Result<CacheManagerRef> {
486    let write_cache = write_cache_from_config(config, puffin_manager, intermediate_manager, None)
487        .await
488        .map_err(|e| {
489            error::IllegalConfigSnafu {
490                msg: format!("Failed to build write cache: {e:?}"),
491            }
492            .build()
493        })?;
494    let cache_manager = Arc::new(
495        CacheManager::builder()
496            .sst_meta_cache_size(config.sst_meta_cache_size.as_bytes())
497            .vector_cache_size(config.vector_cache_size.as_bytes())
498            .page_cache_size(config.page_cache_size.as_bytes())
499            .selector_result_cache_size(config.selector_result_cache_size.as_bytes())
500            .range_result_cache_size(config.range_result_cache_size.as_bytes())
501            .prefilter_result_cache_size(config.prefilter_result_cache_size.as_bytes())
502            .index_metadata_size(config.index.metadata_cache_size.as_bytes())
503            .index_content_size(config.index.content_cache_size.as_bytes())
504            .index_content_page_size(config.index.content_cache_page_size.as_bytes())
505            .index_result_cache_size(config.index.result_cache_size.as_bytes())
506            .puffin_metadata_size(config.index.metadata_cache_size.as_bytes())
507            .write_cache(write_cache)
508            .build(),
509    );
510    Ok(cache_manager)
511}
512
513fn new_noop_file_purger() -> FilePurgerRef {
514    #[derive(Debug)]
515    struct Noop;
516    impl FilePurger for Noop {
517        fn remove_file(&self, _file_meta: FileMeta, _is_delete: bool, _index_outdated: bool) {}
518    }
519    Arc::new(Noop)
520}
521
522async fn load_parquet_metadata(
523    object_store: ObjectStore,
524    path: &str,
525    file_size: u64,
526) -> Result<parquet::file::metadata::ParquetMetaData, Box<dyn std::error::Error + Send + Sync>> {
527    use parquet::file::FOOTER_SIZE;
528    use parquet::file::metadata::ParquetMetaDataReader;
529    let actual_size = if file_size == 0 {
530        object_store.stat(path).await?.content_length()
531    } else {
532        file_size
533    };
534    if actual_size < FOOTER_SIZE as u64 {
535        return Err("file too small".into());
536    }
537    let prefetch: u64 = 64 * 1024;
538    let start = actual_size.saturating_sub(prefetch);
539    let buffer = object_store
540        .read_with(path)
541        .range(start..actual_size)
542        .await?
543        .to_vec();
544    let buffer_len = buffer.len();
545    let mut footer = [0; 8];
546    footer.copy_from_slice(&buffer[buffer_len - FOOTER_SIZE..]);
547    let footer = FooterTail::try_new(&footer)?;
548    let metadata_len = footer.metadata_length() as u64;
549    if actual_size - (FOOTER_SIZE as u64) < metadata_len {
550        return Err("invalid footer/metadata length".into());
551    }
552    if (metadata_len as usize) <= buffer_len - FOOTER_SIZE {
553        let metadata_start = buffer_len - metadata_len as usize - FOOTER_SIZE;
554        let meta = ParquetMetaDataReader::decode_metadata(
555            &buffer[metadata_start..buffer_len - FOOTER_SIZE],
556        )?;
557        Ok(meta)
558    } else {
559        let metadata_start = actual_size - metadata_len - FOOTER_SIZE as u64;
560        let data = object_store
561            .read_with(path)
562            .range(metadata_start..(actual_size - FOOTER_SIZE as u64))
563            .await?
564            .to_vec();
565        let meta = ParquetMetaDataReader::decode_metadata(&data)?;
566        Ok(meta)
567    }
568}
569
570#[cfg(test)]
571mod tests {
572    use std::path::PathBuf;
573    use std::str::FromStr;
574
575    use common_base::readable_size::ReadableSize;
576    use store_api::region_request::PathType;
577
578    use crate::datanode::objbench::parse_file_dir_components;
579    use crate::datanode::tool_util::parse_config;
580
581    #[test]
582    fn test_parse_dir() {
583        let meta_path = "data/greptime/public/1024/1024_0000000000/metadata/00020380-009c-426d-953e-b4e34c15af34.parquet";
584        let c = parse_file_dir_components(meta_path).unwrap();
585        assert_eq!(
586            c.file_id.to_string(),
587            "00020380-009c-426d-953e-b4e34c15af34"
588        );
589        assert_eq!(c.catalog, "greptime");
590        assert_eq!(c.schema, "public");
591        assert_eq!(c.table_id, 1024);
592        assert_eq!(c.region_sequence, 0);
593        assert_eq!(c.path_type, PathType::Metadata);
594
595        let c = parse_file_dir_components(
596            "data/greptime/public/1024/1024_0000000000/data/00020380-009c-426d-953e-b4e34c15af34.parquet",
597        ).unwrap();
598        assert_eq!(
599            c.file_id.to_string(),
600            "00020380-009c-426d-953e-b4e34c15af34"
601        );
602        assert_eq!(c.catalog, "greptime");
603        assert_eq!(c.schema, "public");
604        assert_eq!(c.table_id, 1024);
605        assert_eq!(c.region_sequence, 0);
606        assert_eq!(c.path_type, PathType::Data);
607
608        let c = parse_file_dir_components(
609            "data/greptime/public/1024/1024_0000000000/00020380-009c-426d-953e-b4e34c15af34.parquet",
610        ).unwrap();
611        assert_eq!(
612            c.file_id.to_string(),
613            "00020380-009c-426d-953e-b4e34c15af34"
614        );
615        assert_eq!(c.catalog, "greptime");
616        assert_eq!(c.schema, "public");
617        assert_eq!(c.table_id, 1024);
618        assert_eq!(c.region_sequence, 0);
619        assert_eq!(c.path_type, PathType::Bare);
620    }
621
622    #[test]
623    fn test_parse_config() {
624        let path = "../../config/datanode.example.toml";
625        let (storage, engine, _wal) = parse_config(&PathBuf::from_str(path).unwrap()).unwrap();
626        assert_eq!(storage.data_home, "./greptimedb_data");
627        assert_eq!(engine.index.staging_size, ReadableSize::gb(2));
628    }
629}