1use 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#[derive(Debug, Parser)]
51pub struct ObjbenchCommand {
52 #[clap(long, value_name = "FILE")]
54 pub config: PathBuf,
55
56 #[clap(long, value_name = "PATH")]
58 pub source: String,
59
60 #[clap(short, long, default_value_t = false)]
62 pub verbose: bool,
63
64 #[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 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 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 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 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 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 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 println!("{}", "Writing SST...".yellow());
218
219 #[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 #[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 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 println!(" {}: {:?}", "Metrics".bold(), metrics,);
320
321 println!(" {}: {:?}", "Index".bold(), infos[0].index_metadata);
323
324 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 let pattern =
396 r"^data/([^/]+)/([^/]+)/([^/]+)/([^/]+)_([^/]+)(?:/data|/metadata)?/(.+).parquet$";
397
398 let re = Regex::new(pattern).expect("Invalid regex pattern");
400
401 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 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 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}