Skip to main content

cmd/datanode/
parquetbench.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::{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/// Parquet benchmark command - benchmarks scanning a single parquet SST directly.
54#[derive(Debug, Parser)]
55pub struct ParquetbenchCommand {
56    /// Path to config TOML file (same format as standalone/datanode config)
57    #[clap(long, value_name = "FILE")]
58    config: Option<PathBuf>,
59
60    /// Region ID: either numeric u64 (e.g. "4398046511104") or "table_id:region_num" (e.g. "1024:0")
61    #[clap(long)]
62    region_id: Option<String>,
63
64    /// Table directory relative to data home (e.g. "data/greptime/public/1024/")
65    #[clap(long)]
66    table_dir: Option<String>,
67
68    /// SST file id to benchmark.
69    #[clap(long)]
70    file_id: Option<String>,
71
72    /// Local parquet SST file to benchmark with the direct reader.
73    #[clap(long, value_name = "FILE")]
74    file_path: Option<PathBuf>,
75
76    /// Path to scan request JSON config file (supports projection_names only)
77    #[clap(long, value_name = "FILE")]
78    scan_config: Option<PathBuf>,
79
80    /// Number of iterations for benchmarking
81    #[clap(long, default_value = "1")]
82    iterations: usize,
83
84    /// Number of rows per record batch for the direct reader.
85    #[clap(long, default_value_t = DEFAULT_READ_BATCH_SIZE, value_parser = parse_batch_size)]
86    batch_size: usize,
87
88    /// Path type for the region: bare, data, metadata
89    #[clap(long, default_value = "bare")]
90    path_type: String,
91
92    /// Verbose output
93    #[clap(short, long, default_value_t = false)]
94    verbose: bool,
95
96    /// Output pprof flamegraph
97    #[clap(long, value_name = "FILE")]
98    pprof_file: Option<PathBuf>,
99
100    /// Start pprof after the first iteration (use first iteration as warmup).
101    #[clap(long, default_value_t = false)]
102    pprof_after_warmup: bool,
103
104    /// Read the `__primary_key` column as BinaryArray instead of DictionaryArray.
105    /// Only affects the `direct` reader.
106    #[clap(long, default_value_t = false)]
107    pk_as_binary: bool,
108
109    /// Reader implementation to benchmark.
110    #[clap(long, value_enum, default_value = "direct")]
111    reader: ReaderMode,
112}
113
114#[derive(Debug, Clone, Copy, PartialEq, Eq, ValueEnum)]
115enum ReaderMode {
116    /// Read directly via the push-decoder parquet stream.
117    Direct,
118    /// Read via ParquetReaderBuilder, FileRange, and FlatPruneReader.
119    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    /// Number of columns in the output record batches.
133    columns: usize,
134    /// Schema of the output record batches, captured from the first batch.
135    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, &region_meta)?
218        } else {
219            None
220        };
221        let projection_column_ids = if self.reader == ReaderMode::FlatPrune {
222            resolve_projection_column_ids(&scan_config, &region_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            &region_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, &region_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
928/// Sums the compressed byte size of the selected row groups, used as the denominator for
929/// scan throughput so partial-row-group benchmarks aren't measured against the whole file.
930fn 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}