1use std::path::{Path, PathBuf};
16use std::sync::Arc;
17use std::time::{Duration, Instant};
18
19use clap::{Parser, ValueEnum};
20use colored::Colorize;
21use datatypes::arrow::datatypes::{DataType as ArrowDataType, Field, Schema, SchemaRef};
22use futures::StreamExt;
23use mito2::cache::CacheStrategy;
24use mito2::read::range::FileRangeBuilder;
25use mito2::read::read_columns::ReadColumns;
26use mito2::sst::file::{FileHandle, FileMeta, RegionFileId};
27use mito2::sst::file_purger::NoopFilePurger;
28use mito2::sst::location::sst_file_path;
29use mito2::sst::parquet::metadata::MetadataLoader;
30use mito2::sst::parquet::push_decoder::{
31 SstParquetRangeFetcher, build_sst_parquet_record_batch_stream,
32};
33use mito2::sst::parquet::reader::{MetadataCacheMetrics, ParquetReaderBuilder, ReaderMetrics};
34use mito2::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema};
35use parquet::arrow::ProjectionMask;
36use parquet::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions};
37use serde::Deserialize;
38use smallvec::SmallVec;
39use snafu::ResultExt;
40use store_api::metadata::{RegionMetadata, RegionMetadataRef};
41use store_api::region_request::PathType;
42use store_api::storage::consts::{PRIMARY_KEY_COLUMN_NAME, is_internal_column};
43use store_api::storage::{ColumnId, FileId, RegionId};
44
45use crate::datanode::tool_util::{
46 build_object_store, extract_region_metadata, format_bytes, parse_config, parse_file_id,
47 parse_path_type, parse_region_id,
48};
49use crate::error;
50
51const DEFAULT_READ_BATCH_SIZE: usize = 8 * 1024;
52
53#[derive(Debug, Parser)]
55pub struct ParquetbenchCommand {
56 #[clap(long, value_name = "FILE")]
58 config: Option<PathBuf>,
59
60 #[clap(long)]
62 region_id: Option<String>,
63
64 #[clap(long)]
66 table_dir: Option<String>,
67
68 #[clap(long)]
70 file_id: Option<String>,
71
72 #[clap(long, value_name = "FILE")]
74 file_path: Option<PathBuf>,
75
76 #[clap(long, value_name = "FILE")]
78 scan_config: Option<PathBuf>,
79
80 #[clap(long, default_value = "1")]
82 iterations: usize,
83
84 #[clap(long, default_value_t = DEFAULT_READ_BATCH_SIZE, value_parser = parse_batch_size)]
86 batch_size: usize,
87
88 #[clap(long, default_value = "bare")]
90 path_type: String,
91
92 #[clap(short, long, default_value_t = false)]
94 verbose: bool,
95
96 #[clap(long, value_name = "FILE")]
98 pprof_file: Option<PathBuf>,
99
100 #[clap(long, default_value_t = false)]
102 pprof_after_warmup: bool,
103
104 #[clap(long, default_value_t = false)]
107 pk_as_binary: bool,
108
109 #[clap(long, value_enum, default_value = "direct")]
111 reader: ReaderMode,
112}
113
114#[derive(Debug, Clone, Copy, PartialEq, Eq, ValueEnum)]
115enum ReaderMode {
116 Direct,
118 FlatPrune,
120}
121
122#[derive(Debug, Deserialize, Default)]
123struct ParquetScanConfig {
124 projection_names: Option<Vec<String>>,
125 row_groups: Option<Vec<usize>>,
126}
127
128#[derive(Debug, Default, Clone)]
129struct IterationStats {
130 rows: usize,
131 record_batches: usize,
132 columns: usize,
134 schema: Option<SchemaRef>,
136 elapsed: Duration,
137}
138
139#[derive(Debug)]
140enum ParquetbenchInput {
141 LocalFile {
142 file_path: PathBuf,
143 },
144 Region {
145 config: PathBuf,
146 region_id: String,
147 table_dir: String,
148 file_id: String,
149 path_type: PathType,
150 },
151}
152
153struct ParquetbenchSource {
154 object_store: object_store::ObjectStore,
155 file_path: String,
156 display_path: String,
157 region_id_label: String,
158 region_id: RegionId,
159 file_id: FileId,
160 region_file_id: RegionFileId,
161 path_type: Option<PathType>,
162 table_dir: Option<String>,
163}
164
165impl ParquetbenchCommand {
166 pub async fn run(&self) -> error::Result<()> {
167 if self.verbose {
168 common_telemetry::init_default_ut_logging();
169 }
170
171 println!("{}", "Starting parquetbench...".cyan().bold());
172
173 if self.iterations <= 1 && self.pprof_after_warmup && self.pprof_file.is_some() {
174 return error::IllegalConfigSnafu {
175 msg: "pprof-after-warmup requires at least 2 iterations (1 warmup + 1 profiled)"
176 .to_string(),
177 }
178 .fail();
179 }
180
181 let input = self.resolve_input()?;
182 let mut source = build_source(input).await?;
183
184 let file_size = source
185 .object_store
186 .stat(&source.file_path)
187 .await
188 .map_err(|e| {
189 error::IllegalConfigSnafu {
190 msg: format!("stat failed for {}: {}", source.display_path, e),
191 }
192 .build()
193 })?
194 .content_length();
195 let mut metadata_metrics = MetadataCacheMetrics::default();
196 let parquet_meta =
197 MetadataLoader::new(source.object_store.clone(), &source.file_path, file_size)
198 .load(&mut metadata_metrics)
199 .await
200 .map_err(|e| {
201 error::IllegalConfigSnafu {
202 msg: format!(
203 "read parquet metadata failed for {}: {:?}",
204 source.display_path, e
205 ),
206 }
207 .build()
208 })?;
209 let region_meta = extract_region_metadata(&source.display_path, &parquet_meta)?;
210 if source.table_dir.is_none() {
211 source.region_id = region_meta.region_id;
212 source.region_id_label = source.region_id.as_u64().to_string();
213 source.region_file_id = RegionFileId::new(source.region_id, source.file_id);
214 }
215 let scan_config = self.load_scan_config().await?;
216 let projection = if self.reader == ReaderMode::Direct {
217 resolve_projection_names(&scan_config, ®ion_meta)?
218 } else {
219 None
220 };
221 let projection_column_ids = if self.reader == ReaderMode::FlatPrune {
222 resolve_projection_column_ids(&scan_config, ®ion_meta)?
223 } else {
224 None
225 };
226 let row_groups = resolve_row_groups(&scan_config, parquet_meta.num_row_groups())?;
227 let read_all_row_groups = scan_config.row_groups.is_none();
228 let scanned_bytes = scanned_row_group_bytes(&parquet_meta, &row_groups);
229 let projected_columns = projection_names_display(&scan_config);
230 let row_groups_display = row_groups_display(&row_groups, parquet_meta.num_row_groups());
231 let mut sst_schema = to_flat_sst_arrow_schema(
232 ®ion_meta,
233 &FlatSchemaOptions::from_encoding(region_meta.primary_key_encoding),
234 );
235 if self.pk_as_binary {
236 sst_schema = override_pk_to_binary(&sst_schema);
237 }
238
239 println!(
240 "{} Reader: {}",
241 "✓".green(),
242 match self.reader {
243 ReaderMode::Direct => "direct",
244 ReaderMode::FlatPrune => "flat-prune",
245 }
246 .cyan()
247 );
248 println!(
249 "{} Region ID: {} (u64: {})",
250 "✓".green(),
251 source.region_id_label,
252 source.region_id.as_u64()
253 );
254 println!("{} File path: {}", "✓".green(), source.display_path.cyan());
255 println!(
256 "{} Columns: {}",
257 "✓".green(),
258 projected_columns.as_deref().unwrap_or("all columns").cyan()
259 );
260 println!("{} Row groups: {}", "✓".green(), row_groups_display.cyan());
261 if !read_all_row_groups {
262 println!(
263 "{} Scanned bytes (selected row groups): {}",
264 "✓".green(),
265 format_bytes(scanned_bytes).cyan()
266 );
267 }
268 println!(
269 "{} __primary_key type: {}",
270 "✓".green(),
271 if self.pk_as_binary {
272 "Binary"
273 } else {
274 "Dictionary(UInt32, Binary)"
275 }
276 .cyan()
277 );
278 match self.reader {
279 ReaderMode::Direct => {
280 println!("{} Batch size: {}", "✓".green(), self.batch_size);
281 }
282 ReaderMode::FlatPrune => {
283 println!(
284 "{} Batch size: {} (flat-prune internal default)",
285 "✓".green(),
286 DEFAULT_READ_BATCH_SIZE
287 );
288 println!(
289 "{} --batch-size is only used by the direct reader; ignoring {} in flat-prune mode",
290 "ℹ".blue(),
291 self.batch_size
292 );
293 }
294 }
295 println!(
296 "{} Parquet rows: {}, row groups: {}, file size: {}",
297 "✓".green(),
298 parquet_meta.file_metadata().num_rows(),
299 parquet_meta.num_row_groups(),
300 format_bytes(file_size)
301 );
302 println!(
303 "{} Metadata reads: {}, bytes: {}",
304 "✓".green(),
305 metadata_metrics.num_reads,
306 format_bytes(metadata_metrics.bytes_read)
307 );
308
309 #[cfg(unix)]
310 let mut profiler_guard = if self.pprof_file.is_some() && !self.pprof_after_warmup {
311 println!("{} Starting profiling...", "⚡".yellow());
312 Some(
313 pprof::ProfilerGuardBuilder::default()
314 .frequency(99)
315 .blocklist(&["libc", "libgcc", "pthread", "vdso"])
316 .build()
317 .map_err(|e| {
318 error::IllegalConfigSnafu {
319 msg: format!("Failed to start profiler: {e}"),
320 }
321 .build()
322 })?,
323 )
324 } else {
325 None
326 };
327
328 #[cfg(not(unix))]
329 if self.pprof_file.is_some() {
330 eprintln!(
331 "{}: Profiling is not supported on this platform",
332 "Warning".yellow()
333 );
334 }
335
336 let mut total_elapsed_all = Duration::ZERO;
337 let mut total_rows_all = 0usize;
338 let mut total_batches_all = 0usize;
339 let mut schema_printed = false;
340 let file_handle = FileHandle::new(
341 FileMeta {
342 region_id: source.region_id,
343 file_id: source.file_id,
344 time_range: Default::default(),
345 level: 0,
346 file_size,
347 max_row_group_uncompressed_size: 0,
348 available_indexes: Default::default(),
349 indexes: Default::default(),
350 index_file_size: 0,
351 index_version: 0,
352 num_rows: parquet_meta.file_metadata().num_rows() as u64,
353 num_row_groups: parquet_meta.num_row_groups() as u64,
354 sequence: None,
355 partition_expr: None,
356 num_series: 0,
357 primary_key_min: None,
358 primary_key_max: None,
359 preserve_row_sequence: false,
360 },
361 Arc::new(NoopFilePurger),
362 );
363
364 for iteration in 0..self.iterations {
365 let stats = match self.reader {
366 ReaderMode::Direct => {
367 run_direct_iteration(
368 source.object_store.clone(),
369 source.file_path.clone(),
370 source.region_file_id,
371 parquet_meta.clone(),
372 projection.clone(),
373 row_groups.clone(),
374 sst_schema.clone(),
375 self.batch_size,
376 )
377 .await?
378 }
379 ReaderMode::FlatPrune => {
380 run_flat_prune_iteration(
381 source.object_store.clone(),
382 source.table_dir.clone().ok_or_else(|| {
383 error::IllegalConfigSnafu {
384 msg: "flat-prune reader requires --table-dir".to_string(),
385 }
386 .build()
387 })?,
388 source.path_type.ok_or_else(|| {
389 error::IllegalConfigSnafu {
390 msg: "flat-prune reader requires --path-type".to_string(),
391 }
392 .build()
393 })?,
394 file_handle.clone(),
395 region_meta.clone(),
396 projection_column_ids.clone(),
397 row_groups.clone(),
398 read_all_row_groups,
399 )
400 .await?
401 }
402 };
403
404 total_elapsed_all += stats.elapsed;
405 total_rows_all += stats.rows;
406 total_batches_all += stats.record_batches;
407
408 if !schema_printed && let Some(schema) = &stats.schema {
409 println!(
410 "{} Output schema ({} columns):",
411 "✓".green(),
412 schema.fields().len()
413 );
414 for field in schema.fields() {
415 println!(" - {}: {}", field.name().cyan(), field.data_type());
416 }
417 schema_printed = true;
418 }
419
420 println!(
421 " Iteration {}: {} rows, {} columns, {} record batches in {:?} ({}/s, {}/s)",
422 iteration + 1,
423 stats.rows,
424 stats.columns,
425 stats.record_batches,
426 stats.elapsed,
427 format_rate(stats.rows as f64 / stats.elapsed.as_secs_f64()),
428 format_bytes_per_sec(scanned_bytes as f64 / stats.elapsed.as_secs_f64()),
429 );
430
431 #[cfg(unix)]
432 if iteration == 0 && self.pprof_after_warmup && self.pprof_file.is_some() {
433 println!("{} Starting profiling after warmup...", "⚡".yellow());
434 profiler_guard = Some(
435 pprof::ProfilerGuardBuilder::default()
436 .frequency(99)
437 .blocklist(&["libc", "libgcc", "pthread", "vdso"])
438 .build()
439 .map_err(|e| {
440 error::IllegalConfigSnafu {
441 msg: format!("Failed to start profiler: {e}"),
442 }
443 .build()
444 })?,
445 );
446 }
447 }
448
449 #[cfg(unix)]
450 if let (Some(guard), Some(pprof_file)) = (profiler_guard, &self.pprof_file) {
451 println!("{} Generating flamegraph...", "🔥".yellow());
452 match guard.report().build() {
453 Ok(report) => {
454 let mut flamegraph_data = Vec::new();
455 if let Err(e) = report.flamegraph(&mut flamegraph_data) {
456 println!("{}: Failed to generate flamegraph: {}", "Error".red(), e);
457 } else if let Err(e) = std::fs::write(pprof_file, flamegraph_data) {
458 println!(
459 "{}: Failed to write flamegraph to {}: {}",
460 "Error".red(),
461 pprof_file.display(),
462 e
463 );
464 } else {
465 println!(
466 "{} Flamegraph saved to {}",
467 "✓".green(),
468 pprof_file.display().to_string().cyan()
469 );
470 }
471 }
472 Err(e) => {
473 println!("{}: Failed to generate pprof report: {}", "Error".red(), e);
474 }
475 }
476 }
477
478 if self.iterations > 1 {
479 let avg_elapsed = total_elapsed_all / self.iterations as u32;
480 let avg_rows = total_rows_all / self.iterations;
481 let avg_batches = total_batches_all / self.iterations;
482 println!(
483 "\n{} Average: {} rows, {} record batches in {:?} over {} iterations",
484 "ℹ".blue(),
485 avg_rows,
486 avg_batches,
487 avg_elapsed,
488 self.iterations
489 );
490 }
491
492 println!("\n{}", "Benchmark completed!".green().bold());
493 Ok(())
494 }
495
496 fn resolve_input(&self) -> error::Result<ParquetbenchInput> {
497 let has_region_args = self.config.is_some()
498 || self.region_id.is_some()
499 || self.table_dir.is_some()
500 || self.file_id.is_some();
501
502 if let Some(file_path) = &self.file_path {
503 if self.reader == ReaderMode::FlatPrune {
504 return Err(error::IllegalConfigSnafu {
505 msg: "--file-path currently supports only --reader direct".to_string(),
506 }
507 .build());
508 }
509 if has_region_args {
510 return Err(error::IllegalConfigSnafu {
511 msg: "--file-path cannot be used with --config, --region-id, --table-dir, or --file-id".to_string(),
512 }
513 .build());
514 }
515 return Ok(ParquetbenchInput::LocalFile {
516 file_path: file_path.clone(),
517 });
518 }
519
520 let config = self.config.clone().ok_or_else(|| {
521 error::IllegalConfigSnafu {
522 msg: "missing --config unless --file-path is specified".to_string(),
523 }
524 .build()
525 })?;
526 let region_id = self.region_id.clone().ok_or_else(|| {
527 error::IllegalConfigSnafu {
528 msg: "missing --region-id unless --file-path is specified".to_string(),
529 }
530 .build()
531 })?;
532 let table_dir = self.table_dir.clone().ok_or_else(|| {
533 error::IllegalConfigSnafu {
534 msg: "missing --table-dir unless --file-path is specified".to_string(),
535 }
536 .build()
537 })?;
538 let file_id = self.file_id.clone().ok_or_else(|| {
539 error::IllegalConfigSnafu {
540 msg: "missing --file-id unless --file-path is specified".to_string(),
541 }
542 .build()
543 })?;
544 let path_type = parse_path_type(&self.path_type)?;
545
546 Ok(ParquetbenchInput::Region {
547 config,
548 region_id,
549 table_dir,
550 file_id,
551 path_type,
552 })
553 }
554
555 async fn load_scan_config(&self) -> error::Result<ParquetScanConfig> {
556 if let Some(path) = &self.scan_config {
557 let content = tokio::fs::read_to_string(path)
558 .await
559 .context(error::FileIoSnafu)?;
560 serde_json::from_str::<ParquetScanConfig>(&content).context(error::SerdeJsonSnafu)
561 } else {
562 Ok(ParquetScanConfig::default())
563 }
564 }
565}
566
567async fn build_source(input: ParquetbenchInput) -> error::Result<ParquetbenchSource> {
568 match input {
569 ParquetbenchInput::LocalFile { file_path } => build_local_file_source(&file_path),
570 ParquetbenchInput::Region {
571 config,
572 region_id,
573 table_dir,
574 file_id,
575 path_type,
576 } => build_region_source(&config, ®ion_id, table_dir, &file_id, path_type).await,
577 }
578}
579
580fn build_local_file_source(file_path: &Path) -> error::Result<ParquetbenchSource> {
581 let file_path = std::fs::canonicalize(file_path).map_err(|e| {
582 error::IllegalConfigSnafu {
583 msg: format!("invalid --file-path {}: {e}", file_path.display()),
584 }
585 .build()
586 })?;
587 if !file_path.is_file() {
588 return Err(error::IllegalConfigSnafu {
589 msg: format!("--file-path {} is not a file", file_path.display()),
590 }
591 .build());
592 }
593 let parent = file_path.parent().ok_or_else(|| {
594 error::IllegalConfigSnafu {
595 msg: format!(
596 "--file-path {} has no parent directory",
597 file_path.display()
598 ),
599 }
600 .build()
601 })?;
602 let file_name = file_path
603 .file_name()
604 .and_then(|name| name.to_str())
605 .ok_or_else(|| {
606 error::IllegalConfigSnafu {
607 msg: format!("invalid UTF-8 file name in {}", file_path.display()),
608 }
609 .build()
610 })?
611 .to_string();
612 let object_store = object_store::ObjectStore::new(object_store::services::Fs::default().root(
613 parent.to_str().ok_or_else(|| {
614 error::IllegalConfigSnafu {
615 msg: format!("invalid UTF-8 parent directory in {}", file_path.display()),
616 }
617 .build()
618 })?,
619 ))
620 .map_err(|e| {
621 error::IllegalConfigSnafu {
622 msg: format!("failed to build local file object store: {e:?}"),
623 }
624 .build()
625 })?;
626
627 let file_id = file_path
628 .file_stem()
629 .and_then(|stem| stem.to_str())
630 .and_then(|stem| FileId::parse_str(stem).ok())
631 .unwrap_or_else(FileId::random);
632 let region_id = RegionId::new(0, 0);
633 let region_file_id = RegionFileId::new(region_id, file_id);
634 let display_path = file_path.display().to_string();
635
636 Ok(ParquetbenchSource {
637 object_store,
638 file_path: file_name,
639 display_path,
640 region_id_label: region_id.as_u64().to_string(),
641 region_id,
642 file_id,
643 region_file_id,
644 path_type: None,
645 table_dir: None,
646 })
647}
648
649async fn build_region_source(
650 config: &Path,
651 region_id: &str,
652 table_dir: String,
653 file_id: &str,
654 path_type: PathType,
655) -> error::Result<ParquetbenchSource> {
656 let region = parse_region_id(region_id)?;
657 let file_id = parse_file_id(file_id)?;
658 let region_file_id = RegionFileId::new(region, file_id);
659 let file_path = sst_file_path(&table_dir, region_file_id, path_type);
660
661 let (store_cfg, _mito_config, _wal_config) = parse_config(config)?;
662 let object_store = build_object_store(&store_cfg).await?;
663
664 Ok(ParquetbenchSource {
665 object_store,
666 display_path: file_path.clone(),
667 file_path,
668 region_id_label: region_id.to_string(),
669 region_id: region,
670 file_id,
671 region_file_id,
672 path_type: Some(path_type),
673 table_dir: Some(table_dir),
674 })
675}
676
677#[allow(clippy::too_many_arguments)]
678async fn run_direct_iteration(
679 object_store: object_store::ObjectStore,
680 file_path: String,
681 region_file_id: RegionFileId,
682 parquet_meta: parquet::file::metadata::ParquetMetaData,
683 projection: Option<Vec<usize>>,
684 row_groups: Vec<usize>,
685 sst_schema: SchemaRef,
686 batch_size: usize,
687) -> error::Result<IterationStats> {
688 let parquet_meta = Arc::new(parquet_meta);
689 let arrow_metadata = ArrowReaderMetadata::try_new(
690 parquet_meta.clone(),
691 ArrowReaderOptions::new().with_schema(sst_schema),
692 )
693 .map_err(|e| {
694 error::IllegalConfigSnafu {
695 msg: format!(
696 "Failed to build parquet arrow metadata for {}: {}",
697 file_path, e
698 ),
699 }
700 .build()
701 })?;
702 let projection_mask = match projection.as_ref() {
703 Some(projection) => {
704 ProjectionMask::roots(arrow_metadata.parquet_schema(), projection.iter().copied())
705 }
706 None => ProjectionMask::all(),
707 };
708 let start = Instant::now();
709 let mut stats = IterationStats::default();
710 for row_group_idx in row_groups {
711 let fetcher = SstParquetRangeFetcher::new(
712 region_file_id,
713 file_path.clone(),
714 object_store.clone(),
715 CacheStrategy::Disabled,
716 row_group_idx,
717 None,
718 );
719 let mut stream = build_sst_parquet_record_batch_stream(
720 arrow_metadata.clone(),
721 row_group_idx,
722 None,
723 projection_mask.clone(),
724 fetcher,
725 file_path.clone(),
726 batch_size,
727 )
728 .map_err(|e| {
729 error::IllegalConfigSnafu {
730 msg: format!(
731 "Failed to build parquet record batch stream for {}: {e:?}",
732 file_path
733 ),
734 }
735 .build()
736 })?;
737 while let Some(batch) = stream.next().await.transpose().map_err(|e| {
738 error::IllegalConfigSnafu {
739 msg: format!("Failed to scan parquet file {}: {e:?}", file_path),
740 }
741 .build()
742 })? {
743 stats.rows += batch.num_rows();
744 stats.record_batches += 1;
745 stats.columns = batch.num_columns();
746 if stats.schema.is_none() {
747 stats.schema = Some(batch.schema());
748 }
749 }
750 }
751 stats.elapsed = start.elapsed();
752 Ok(stats)
753}
754
755#[allow(clippy::too_many_arguments)]
756async fn run_flat_prune_iteration(
757 object_store: object_store::ObjectStore,
758 table_dir: String,
759 path_type: PathType,
760 file_handle: FileHandle,
761 region_meta: RegionMetadataRef,
762 projection: Option<Vec<ColumnId>>,
763 row_groups: Vec<usize>,
764 read_all_row_groups: bool,
765) -> error::Result<IterationStats> {
766 let reader_builder = ParquetReaderBuilder::new(table_dir, path_type, file_handle, object_store)
767 .expected_metadata(Some(region_meta))
768 .cache(CacheStrategy::Disabled)
769 .projection(projection.map(ReadColumns::new));
770 let mut reader_metrics = ReaderMetrics::default();
771 let start = Instant::now();
772 let mut stats = IterationStats::default();
773 let Some((context, selection)) = reader_builder
774 .build_reader_input(&mut reader_metrics)
775 .await
776 .map_err(|e| {
777 error::IllegalConfigSnafu {
778 msg: format!("build flat prune reader input failed: {e:?}"),
779 }
780 .build()
781 })?
782 else {
783 stats.elapsed = start.elapsed();
784 return Ok(stats);
785 };
786
787 let range_builder = FileRangeBuilder::new(Arc::new(context), selection);
788 let mut ranges = SmallVec::new();
789 if read_all_row_groups {
790 range_builder.build_ranges(-1, &mut ranges);
791 } else {
792 for row_group_idx in row_groups {
793 range_builder.build_ranges(row_group_idx as i64, &mut ranges);
794 }
795 }
796
797 for range in ranges {
798 let Some(mut reader) = range.flat_reader(None, None).await.map_err(|e| {
799 error::IllegalConfigSnafu {
800 msg: format!("build flat prune reader failed: {e:?}"),
801 }
802 .build()
803 })?
804 else {
805 continue;
806 };
807 while let Some(batch) = reader.next_batch().await.map_err(|e| {
808 error::IllegalConfigSnafu {
809 msg: format!("scan flat prune reader failed: {e:?}"),
810 }
811 .build()
812 })? {
813 stats.rows += batch.num_rows();
814 stats.record_batches += 1;
815 stats.columns = batch.num_columns();
816 if stats.schema.is_none() {
817 stats.schema = Some(batch.schema());
818 }
819 }
820 }
821
822 stats.elapsed = start.elapsed();
823 Ok(stats)
824}
825
826fn resolve_projection_names(
827 scan_config: &ParquetScanConfig,
828 metadata: &RegionMetadata,
829) -> error::Result<Option<Vec<usize>>> {
830 let Some(projection_names) = &scan_config.projection_names else {
831 return Ok(None);
832 };
833
834 let sst_schema = to_flat_sst_arrow_schema(
835 metadata,
836 &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
837 );
838 let available_columns = sst_schema
839 .fields()
840 .iter()
841 .map(|field| field.name().as_str())
842 .collect::<Vec<_>>()
843 .join(", ");
844 let projection = projection_names
845 .iter()
846 .map(|name| {
847 sst_schema
848 .column_with_name(name)
849 .map(|x| x.0)
850 .ok_or_else(|| {
851 error::IllegalConfigSnafu {
852 msg: format!(
853 "Unknown column '{}' in projection_names, available columns: [{}]",
854 name, available_columns
855 ),
856 }
857 .build()
858 })
859 })
860 .collect::<error::Result<Vec<_>>>()?;
861 Ok(Some(projection))
862}
863
864fn resolve_projection_column_ids(
865 scan_config: &ParquetScanConfig,
866 metadata: &RegionMetadata,
867) -> error::Result<Option<Vec<ColumnId>>> {
868 let Some(projection_names) = &scan_config.projection_names else {
869 return Ok(None);
870 };
871
872 let available_columns = metadata
873 .column_metadatas
874 .iter()
875 .map(|column| column.column_schema.name.as_str())
876 .collect::<Vec<_>>()
877 .join(", ");
878 let projection = projection_names
879 .iter()
880 .filter_map(|name| {
881 if is_internal_column(name) {
882 return None;
883 }
884 Some(
885 metadata
886 .column_metadatas
887 .iter()
888 .find(|column| column.column_schema.name == *name)
889 .map(|column| column.column_id)
890 .ok_or_else(|| {
891 error::IllegalConfigSnafu {
892 msg: format!(
893 "Unknown column '{}' in projection_names, available columns: [{}]",
894 name, available_columns
895 ),
896 }
897 .build()
898 }),
899 )
900 })
901 .collect::<error::Result<Vec<_>>>()?;
902 Ok(Some(projection))
903}
904
905fn resolve_row_groups(
906 scan_config: &ParquetScanConfig,
907 num_row_groups: usize,
908) -> error::Result<Vec<usize>> {
909 match &scan_config.row_groups {
910 Some(row_groups) => {
911 for row_group_idx in row_groups {
912 if *row_group_idx >= num_row_groups {
913 return Err(error::IllegalConfigSnafu {
914 msg: format!(
915 "Invalid row group {} in row_groups, parquet file has row groups [0, {})",
916 row_group_idx, num_row_groups
917 ),
918 }
919 .build());
920 }
921 }
922 Ok(row_groups.clone())
923 }
924 None => Ok((0..num_row_groups).collect()),
925 }
926}
927
928fn scanned_row_group_bytes(
931 parquet_meta: &parquet::file::metadata::ParquetMetaData,
932 row_groups: &[usize],
933) -> u64 {
934 row_groups
935 .iter()
936 .map(|&idx| parquet_meta.row_group(idx).compressed_size() as u64)
937 .sum()
938}
939
940fn projection_names_display(scan_config: &ParquetScanConfig) -> Option<String> {
941 scan_config
942 .projection_names
943 .as_ref()
944 .map(|cols| cols.join(", "))
945}
946
947fn row_groups_display(row_groups: &[usize], total_row_groups: usize) -> String {
948 if row_groups.len() == total_row_groups {
949 "all row groups".to_string()
950 } else {
951 row_groups
952 .iter()
953 .map(|idx| idx.to_string())
954 .collect::<Vec<_>>()
955 .join(", ")
956 }
957}
958
959fn override_pk_to_binary(schema: &SchemaRef) -> SchemaRef {
960 let new_fields: Vec<_> = schema
961 .fields()
962 .iter()
963 .map(|f| {
964 if f.name() == PRIMARY_KEY_COLUMN_NAME {
965 Arc::new(Field::new(
966 PRIMARY_KEY_COLUMN_NAME,
967 ArrowDataType::Binary,
968 f.is_nullable(),
969 ))
970 } else {
971 f.clone()
972 }
973 })
974 .collect();
975 Arc::new(Schema::new(new_fields))
976}
977
978fn parse_batch_size(s: &str) -> Result<usize, String> {
979 let batch_size = s
980 .parse::<usize>()
981 .map_err(|e| format!("invalid batch size '{s}': {e}"))?;
982 if batch_size == 0 {
983 return Err("batch size must be greater than 0".to_string());
984 }
985 Ok(batch_size)
986}
987
988fn format_rate(rate: f64) -> String {
989 if !rate.is_finite() {
990 return "inf rows".to_string();
991 }
992 format!("{rate:.2} rows")
993}
994
995fn format_bytes_per_sec(bytes_per_sec: f64) -> String {
996 if !bytes_per_sec.is_finite() {
997 return "inf B/s".to_string();
998 }
999 format!("{}/s", format_bytes(bytes_per_sec as u64))
1000}
1001
1002#[cfg(test)]
1003mod tests {
1004 use api::v1::SemanticType;
1005 use datatypes::prelude::ConcreteDataType;
1006 use serde_json::json;
1007 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
1008 use store_api::storage::ColumnSchema;
1009
1010 use super::*;
1011
1012 fn test_command() -> ParquetbenchCommand {
1013 ParquetbenchCommand {
1014 config: None,
1015 region_id: None,
1016 table_dir: None,
1017 file_id: None,
1018 file_path: None,
1019 scan_config: None,
1020 iterations: 1,
1021 batch_size: DEFAULT_READ_BATCH_SIZE,
1022 path_type: "bare".to_string(),
1023 verbose: false,
1024 pprof_file: None,
1025 pprof_after_warmup: false,
1026 pk_as_binary: false,
1027 reader: ReaderMode::Direct,
1028 }
1029 }
1030
1031 fn new_test_metadata() -> RegionMetadata {
1032 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 0));
1033 builder
1034 .push_column_metadata(ColumnMetadata {
1035 column_schema: ColumnSchema::new(
1036 "host",
1037 ConcreteDataType::string_datatype(),
1038 false,
1039 ),
1040 semantic_type: SemanticType::Tag,
1041 column_id: 1,
1042 })
1043 .push_column_metadata(ColumnMetadata {
1044 column_schema: ColumnSchema::new("cpu", ConcreteDataType::float64_datatype(), true),
1045 semantic_type: SemanticType::Field,
1046 column_id: 2,
1047 })
1048 .push_column_metadata(ColumnMetadata {
1049 column_schema: ColumnSchema::new(
1050 "ts",
1051 ConcreteDataType::timestamp_millisecond_datatype(),
1052 false,
1053 ),
1054 semantic_type: SemanticType::Timestamp,
1055 column_id: 3,
1056 })
1057 .primary_key(vec![1]);
1058 builder.build().unwrap()
1059 }
1060
1061 #[test]
1062 fn test_resolve_input_accepts_direct_file_for_direct_reader() {
1063 let mut command = test_command();
1064 command.file_path = Some(PathBuf::from("/tmp/source.parquet"));
1065
1066 let input = command.resolve_input().unwrap();
1067 match input {
1068 ParquetbenchInput::LocalFile { file_path } => {
1069 assert_eq!(file_path, PathBuf::from("/tmp/source.parquet"));
1070 }
1071 ParquetbenchInput::Region { .. } => panic!("expected local file input"),
1072 }
1073 }
1074
1075 #[test]
1076 fn test_resolve_input_rejects_mixed_input_modes() {
1077 let mut command = test_command();
1078 command.file_path = Some(PathBuf::from("/tmp/source.parquet"));
1079 command.config = Some(PathBuf::from("config.toml"));
1080
1081 let err = command.resolve_input().unwrap_err();
1082 assert!(err.to_string().contains("--file-path cannot be used with"));
1083 }
1084
1085 #[test]
1086 fn test_parse_scan_config_projection_names() {
1087 let config: ParquetScanConfig =
1088 serde_json::from_value(json!({ "projection_names": ["host", "ts"] })).unwrap();
1089 assert_eq!(
1090 config.projection_names,
1091 Some(vec!["host".to_string(), "ts".to_string()])
1092 );
1093 }
1094
1095 #[test]
1096 fn test_parse_scan_config_row_groups() {
1097 let config: ParquetScanConfig =
1098 serde_json::from_value(json!({ "row_groups": [0, 2, 4] })).unwrap();
1099 assert_eq!(config.row_groups, Some(vec![0, 2, 4]));
1100 }
1101
1102 #[test]
1103 fn test_resolve_projection_names() {
1104 let metadata = new_test_metadata();
1105 let projection = resolve_projection_names(
1106 &ParquetScanConfig {
1107 projection_names: Some(vec!["cpu".to_string(), "host".to_string()]),
1108 row_groups: None,
1109 },
1110 &metadata,
1111 )
1112 .unwrap();
1113 assert_eq!(projection, Some(vec![1, 0]));
1114 }
1115
1116 #[test]
1117 fn test_resolve_projection_column_ids() {
1118 let metadata = new_test_metadata();
1119 let projection = resolve_projection_column_ids(
1120 &ParquetScanConfig {
1121 projection_names: Some(vec!["cpu".to_string(), "host".to_string()]),
1122 row_groups: None,
1123 },
1124 &metadata,
1125 )
1126 .unwrap();
1127 assert_eq!(projection, Some(vec![2, 1]));
1128 }
1129
1130 #[test]
1131 fn test_resolve_projection_column_ids_ignores_internal_columns() {
1132 let metadata = new_test_metadata();
1133 let projection = resolve_projection_column_ids(
1134 &ParquetScanConfig {
1135 projection_names: Some(vec![
1136 "cpu".to_string(),
1137 "__primary_key".to_string(),
1138 "__sequence".to_string(),
1139 "__op_type".to_string(),
1140 ]),
1141 row_groups: None,
1142 },
1143 &metadata,
1144 )
1145 .unwrap();
1146 assert_eq!(projection, Some(vec![2]));
1147 }
1148
1149 #[test]
1150 fn test_resolve_projection_names_unknown() {
1151 let metadata = new_test_metadata();
1152 let err = resolve_projection_names(
1153 &ParquetScanConfig {
1154 projection_names: Some(vec!["memory".to_string()]),
1155 row_groups: None,
1156 },
1157 &metadata,
1158 )
1159 .unwrap_err();
1160 let msg = err.to_string();
1161 assert!(msg.contains("projection_names"));
1162 assert!(msg.contains("host"));
1163 assert!(msg.contains("cpu"));
1164 assert!(msg.contains("ts"));
1165 }
1166
1167 #[test]
1168 fn test_resolve_row_groups_all() {
1169 assert_eq!(
1170 resolve_row_groups(&ParquetScanConfig::default(), 3).unwrap(),
1171 vec![0, 1, 2]
1172 );
1173 }
1174
1175 #[test]
1176 fn test_resolve_row_groups_subset() {
1177 let config = ParquetScanConfig {
1178 projection_names: None,
1179 row_groups: Some(vec![2, 0]),
1180 };
1181 assert_eq!(resolve_row_groups(&config, 4).unwrap(), vec![2, 0]);
1182 }
1183
1184 #[test]
1185 fn test_resolve_row_groups_invalid() {
1186 let config = ParquetScanConfig {
1187 projection_names: None,
1188 row_groups: Some(vec![3]),
1189 };
1190 let err = resolve_row_groups(&config, 3).unwrap_err();
1191 assert!(err.to_string().contains("Invalid row group 3"));
1192 }
1193
1194 #[test]
1195 fn test_sst_file_path_resolution() {
1196 let file_id = FileId::parse_str("00020380-009c-426d-953e-b4e34c15af34").unwrap();
1197 let region_file_id = RegionFileId::new(RegionId::new(1024, 0), file_id);
1198 assert_eq!(
1199 sst_file_path("data/greptime/public/1024", region_file_id, PathType::Bare),
1200 "data/greptime/public/1024/1024_0000000000/00020380-009c-426d-953e-b4e34c15af34.parquet"
1201 );
1202 }
1203}