diff --git a/.github/workflows/sql-benchmarks.yml b/.github/workflows/sql-benchmarks.yml index 243db83aff8..86e97c59d85 100644 --- a/.github/workflows/sql-benchmarks.yml +++ b/.github/workflows/sql-benchmarks.yml @@ -68,7 +68,6 @@ on: "name": "TPC-H SF=1 on NVME", "data_formats": ["parquet", "vortex", "vortex-compact", "duckdb"], "pr_targets": [ - {"engine": "datafusion", "format": "arrow"}, {"engine": "datafusion", "format": "parquet"}, {"engine": "datafusion", "format": "vortex"}, {"engine": "datafusion", "format": "vortex-compact"}, @@ -78,7 +77,6 @@ on: {"engine": "duckdb", "format": "duckdb"} ], "develop_targets": [ - {"engine": "datafusion", "format": "arrow"}, {"engine": "datafusion", "format": "parquet"}, {"engine": "datafusion", "format": "vortex"}, {"engine": "datafusion", "format": "vortex-compact"}, @@ -123,7 +121,6 @@ on: "name": "TPC-H SF=10 on NVME", "data_formats": ["parquet", "vortex", "vortex-compact", "duckdb"], "pr_targets": [ - {"engine": "datafusion", "format": "arrow"}, {"engine": "datafusion", "format": "parquet"}, {"engine": "datafusion", "format": "vortex"}, {"engine": "datafusion", "format": "vortex-compact"}, @@ -133,7 +130,6 @@ on: {"engine": "duckdb", "format": "duckdb"} ], "develop_targets": [ - {"engine": "datafusion", "format": "arrow"}, {"engine": "datafusion", "format": "parquet"}, {"engine": "datafusion", "format": "vortex"}, {"engine": "datafusion", "format": "vortex-compact"}, @@ -350,14 +346,12 @@ on: "name": "TPC-H SF=1 on NVME", "data_formats": ["parquet", "vortex"], "pr_targets": [ - {"engine": "datafusion", "format": "arrow"}, {"engine": "datafusion", "format": "parquet"}, {"engine": "datafusion", "format": "vortex"}, {"engine": "duckdb", "format": "parquet"}, {"engine": "duckdb", "format": "vortex"} ], "develop_targets": [ - {"engine": "datafusion", "format": "arrow"}, {"engine": "datafusion", "format": "parquet"}, {"engine": "datafusion", "format": "vortex"}, {"engine": "datafusion", "format": "lance"}, @@ -395,14 +389,12 @@ on: "name": "TPC-H SF=10 on NVME", "data_formats": ["parquet", "vortex"], "pr_targets": [ - {"engine": "datafusion", "format": "arrow"}, {"engine": "datafusion", "format": "parquet"}, {"engine": "datafusion", "format": "vortex"}, {"engine": "duckdb", "format": "parquet"}, {"engine": "duckdb", "format": "vortex"} ], "develop_targets": [ - {"engine": "datafusion", "format": "arrow"}, {"engine": "datafusion", "format": "parquet"}, {"engine": "datafusion", "format": "vortex"}, {"engine": "datafusion", "format": "lance"}, diff --git a/bench-orchestrator/README.md b/bench-orchestrator/README.md index 0b267008a85..22b1b1a1249 100644 --- a/bench-orchestrator/README.md +++ b/bench-orchestrator/README.md @@ -269,7 +269,7 @@ vx-bench clean --older-than "30 days" --no-keep-labeled | Engine | Supported Formats | |------------|-------------------------------------------| -| datafusion | arrow, parquet, vortex, vortex-compact, lance | +| datafusion | parquet, vortex, vortex-compact, lance | | duckdb | parquet, vortex, vortex-compact, duckdb | | lance | lance | diff --git a/bench-orchestrator/bench_orchestrator/config.py b/bench-orchestrator/bench_orchestrator/config.py index 4c9df79fef1..20a42f902d1 100644 --- a/bench-orchestrator/bench_orchestrator/config.py +++ b/bench-orchestrator/bench_orchestrator/config.py @@ -31,7 +31,6 @@ def binary_name(self) -> str: class Format(Enum): """Data formats for benchmarks.""" - ARROW = "arrow" PARQUET = "parquet" VORTEX = "vortex" VORTEX_COMPACT = "vortex-compact" @@ -60,7 +59,6 @@ class Benchmark(Enum): # Engine to supported formats mapping. ENGINE_FORMATS: dict[Engine, list[Format]] = { Engine.DATAFUSION: [ - Format.ARROW, Format.PARQUET, Format.VORTEX, Format.VORTEX_COMPACT, diff --git a/bench-orchestrator/tests/test_config.py b/bench-orchestrator/tests/test_config.py index fa3be9a2df6..12c48d21433 100644 --- a/bench-orchestrator/tests/test_config.py +++ b/bench-orchestrator/tests/test_config.py @@ -46,15 +46,15 @@ def test_resolve_axis_targets_offers_vortex_native_on_duckdb_only() -> None: def test_resolve_axis_targets_filters_unsupported_combinations() -> None: targets, warnings = resolve_axis_targets( [Engine.DATAFUSION, Engine.DUCKDB], - [Format.ARROW, Format.PARQUET], + [Format.LANCE, Format.PARQUET], ) assert targets == [ - BenchmarkTarget(engine=Engine.DATAFUSION, format=Format.ARROW), + BenchmarkTarget(engine=Engine.DATAFUSION, format=Format.LANCE), BenchmarkTarget(engine=Engine.DATAFUSION, format=Format.PARQUET), BenchmarkTarget(engine=Engine.DUCKDB, format=Format.PARQUET), ] - assert warnings == ["Format arrow is not supported by engine duckdb"] + assert warnings == ["Format lance is not supported by engine duckdb"] def test_resolve_axis_targets_skips_engines_a_benchmark_cannot_run() -> None: diff --git a/benchmarks/datafusion-bench/src/lib.rs b/benchmarks/datafusion-bench/src/lib.rs index c45aa99be2c..b39120acfce 100644 --- a/benchmarks/datafusion-bench/src/lib.rs +++ b/benchmarks/datafusion-bench/src/lib.rs @@ -7,7 +7,6 @@ pub mod tracer; use std::sync::Arc; use datafusion::datasource::file_format::FileFormat; -use datafusion::datasource::file_format::arrow::ArrowFormat; use datafusion::datasource::file_format::csv::CsvFormat; use datafusion::datasource::file_format::parquet::ParquetFormat; use datafusion::datasource::provider::DefaultTableFactory; @@ -109,7 +108,6 @@ pub fn make_object_store( pub fn format_to_df_format(format: Format) -> Arc { match format { Format::Csv => Arc::new(CsvFormat::default()) as _, - Format::Arrow => Arc::new(ArrowFormat), Format::Parquet => Arc::new(ParquetFormat::new()), Format::OnDiskVortex | Format::VortexCompact | Format::VortexNative => Arc::new( VortexFormat::new_with_options(SESSION.clone(), vortex_table_options()), diff --git a/benchmarks/datafusion-bench/src/main.rs b/benchmarks/datafusion-bench/src/main.rs index 67521a8147e..8c67896c64b 100644 --- a/benchmarks/datafusion-bench/src/main.rs +++ b/benchmarks/datafusion-bench/src/main.rs @@ -15,7 +15,6 @@ use datafusion::datasource::listing::ListingOptions; use datafusion::datasource::listing::ListingTable; use datafusion::datasource::listing::ListingTableConfig; use datafusion::datasource::listing::ListingTableUrl; -use datafusion::parquet::arrow::ParquetRecordBatchStreamBuilder; use datafusion::prelude::SessionContext; use datafusion_bench::format_to_df_format; use datafusion_bench::metrics::MetricsSetExt; @@ -24,10 +23,12 @@ use datafusion_bench::tracer::get_static_tracer; use datafusion_bench::tracer::set_labels; use datafusion_physical_plan::ExecutionPlan; use datafusion_physical_plan::collect; -use futures::StreamExt; use parking_lot::Mutex; -use tokio::fs::File; +use vortex::file::multi::MultiFileDataSource; use vortex::io::filesystem::FileSystemRef; +use vortex::io::object_store::ObjectStoreFileSystem; +use vortex::io::session::RuntimeSessionExt; +use vortex::scan::DataSource as _; use vortex::scan::DataSourceRef; use vortex_arrow::ToArrowType; use vortex_bench::Benchmark; @@ -49,6 +50,7 @@ use vortex_bench::runner::filter_queries; use vortex_bench::setup_logging_and_tracing; use vortex_bench::v3; use vortex_datafusion::metrics::VortexMetricsFinder; +use vortex_datafusion::v2::VortexTable; /// Common arguments shared across benchmarks #[derive(Parser)] @@ -256,43 +258,39 @@ async fn register_benchmark_tables( benchmark: &B, format: Format, ) -> anyhow::Result<()> { - match format { - Format::Arrow => register_arrow_tables(session, benchmark).await, - _ if use_scan_api() && matches!(format, Format::OnDiskVortex | Format::VortexCompact) => { - register_v2_tables(session, benchmark, format).await + if use_scan_api() && matches!(format, Format::OnDiskVortex | Format::VortexCompact) { + register_v2_tables(session, benchmark, format).await + } else { + let benchmark_base = benchmark.data_url().join(&format!("{}/", format.name()))?; + let file_format = format_to_df_format(format); + + for table in benchmark.table_specs().iter() { + let pattern = benchmark.pattern(table.name, format); + let table_url = ListingTableUrl::try_new(benchmark_base.clone(), pattern)?; + + let listing_options = ListingOptions::new(Arc::clone(&file_format)) + .with_session_config_options(session.state().config()); + let mut config = + ListingTableConfig::new(table_url).with_listing_options(listing_options); + + config = match table.schema.as_ref() { + Some(schema) => config.with_schema(Arc::new(schema.clone())), + None => config.infer_schema(&session.state()).await?, + }; + + let listing_table = Arc::new( + ListingTable::try_new(config)?.with_cache( + session + .runtime_env() + .cache_manager + .get_file_statistic_cache(), + ), + ); + + session.register_table(table.name, listing_table)?; } - _ => { - let benchmark_base = benchmark.data_url().join(&format!("{}/", format.name()))?; - let file_format = format_to_df_format(format); - - for table in benchmark.table_specs().iter() { - let pattern = benchmark.pattern(table.name, format); - let table_url = ListingTableUrl::try_new(benchmark_base.clone(), pattern)?; - - let listing_options = ListingOptions::new(Arc::clone(&file_format)) - .with_session_config_options(session.state().config()); - let mut config = - ListingTableConfig::new(table_url).with_listing_options(listing_options); - - config = match table.schema.as_ref() { - Some(schema) => config.with_schema(Arc::new(schema.clone())), - None => config.infer_schema(&session.state()).await?, - }; - - let listing_table = Arc::new( - ListingTable::try_new(config)?.with_cache( - session - .runtime_env() - .cache_manager - .get_file_statistic_cache(), - ), - ); - - session.register_table(table.name, listing_table)?; - } - Ok(()) - } + Ok(()) } } @@ -302,12 +300,6 @@ async fn register_v2_tables( benchmark: &B, format: Format, ) -> anyhow::Result<()> { - use vortex::file::multi::MultiFileDataSource; - use vortex::io::object_store::ObjectStoreFileSystem; - use vortex::io::session::RuntimeSessionExt; - use vortex::scan::DataSource as _; - use vortex_datafusion::v2::VortexTable; - let benchmark_base = benchmark.data_url().join(&format!("{}/", format.name()))?; for table in benchmark.table_specs().iter() { @@ -345,63 +337,6 @@ async fn register_v2_tables( Ok(()) } -/// Load Arrow IPC files into in-memory DataFusion tables. -async fn register_arrow_tables( - session: &SessionContext, - benchmark: &B, -) -> anyhow::Result<()> { - use datafusion::datasource::MemTable; - - let parquet_dir = benchmark - .data_url() - .to_file_path() - .map_err(|_| anyhow::anyhow!("Arrow format requires local file path"))? - .join(Format::Parquet.name()); - - // Read all arrow files from the directory - let data_files = std::fs::read_dir(&parquet_dir)?.collect::, _>>()?; - - for table in benchmark.table_specs().iter() { - let pattern = benchmark.pattern(table.name, Format::Parquet); - - // Find files matching this table's pattern - let matching_files: Vec<_> = data_files - .iter() - .filter(|entry| { - let filename = entry.file_name(); - let filename_str = filename.to_str().unwrap_or(""); - match &pattern { - Some(p) => p.matches(filename_str), - None => filename_str == format!("{}.{}", table.name, Format::Parquet.ext()), - } - }) - .collect(); - - // Load all matching files into memory - let mut all_batches = Vec::new(); - let mut schema = None; - - for dir_entry in matching_files { - let file = File::open(dir_entry.path()).await?; - let mut reader = ParquetRecordBatchStreamBuilder::new(file).await?.build()?; - if schema.is_none() { - schema = Some(reader.schema()).cloned(); - } - - while let Some(batch) = reader.next().await { - all_batches.push(batch?); - } - } - - if let Some(schema) = schema { - let mem_table = MemTable::try_new(schema, vec![all_batches])?; - session.register_table(table.name, Arc::new(mem_table))?; - } - } - - Ok(()) -} - /// Wrapper around DataFusion record batches implementing `BenchmarkQueryResult`. pub struct DataFusionQueryResult(pub Vec); diff --git a/benchmarks/random-access-bench/src/main.rs b/benchmarks/random-access-bench/src/main.rs index 6722265fde6..0a54778ee98 100644 --- a/benchmarks/random-access-bench/src/main.rs +++ b/benchmarks/random-access-bench/src/main.rs @@ -269,7 +269,7 @@ async fn benchmark_random_access( let timing = TimingMeasurement { name: measurement_name.to_string(), storage: storage.to_string(), - target: Target::new(format_to_engine(format), format), + target: Target::new(Engine::default(), format), runs, }; Ok(RandomAccessRun { @@ -326,17 +326,6 @@ fn push_v3_random_access_record(records: &mut Vec, run: &RandomAcc records.push(v3::random_access_record(&run.timing, &dataset)); } -/// Map format to the appropriate engine for random access benchmarks. -fn format_to_engine(format: Format) -> Engine { - match format { - Format::OnDiskVortex | Format::VortexCompact => Engine::Vortex, - Format::Parquet => Engine::Arrow, - #[cfg(feature = "lance")] - Format::Lance => Engine::Arrow, // Is this right here? - _ => Engine::default(), - } -} - /// Open a random accessor for any supported format. /// /// For Vortex and Parquet, the path comes from [`BenchDataset::path`]. @@ -528,7 +517,7 @@ mod tests { RandomAccessRun { timing: TimingMeasurement { name: format!("random-access/{dataset}/parquet-tokio-local-disk"), - target: Target::new(Engine::Arrow, Format::Parquet), + target: Target::new(Engine::Vortex, Format::Parquet), storage: STORAGE_NVME.to_string(), runs: vec![Duration::from_nanos(10)], }, diff --git a/benchmarks/random-access-bench/src/render.rs b/benchmarks/random-access-bench/src/render.rs index a53db41871a..a23d34dc91d 100644 --- a/benchmarks/random-access-bench/src/render.rs +++ b/benchmarks/random-access-bench/src/render.rs @@ -198,7 +198,7 @@ mod tests { RandomAccessRun { timing: TimingMeasurement { name: format!("random-access/{dataset}/{}-tokio-local-disk", format.ext()), - target: Target::new(Engine::Arrow, format), + target: Target::new(Engine::Vortex, format), storage: "nvme".to_string(), runs: vec![Duration::from_micros(micros)], }, @@ -273,10 +273,6 @@ mod tests { rendered.contains("vortex-cached") && rendered.contains("vortex-reopen"), "expected ext-based column headers, got:\n{rendered}" ); - assert!( - !rendered.contains("arrow"), - "expected no engine header row, got:\n{rendered}" - ); assert!( rendered.contains("random-access/taxi") && rendered.contains("random-access/taxi/uniform"), diff --git a/vortex-bench/src/lib.rs b/vortex-bench/src/lib.rs index abf4e1ccea3..bc077a64719 100644 --- a/vortex-bench/src/lib.rs +++ b/vortex-bench/src/lib.rs @@ -141,8 +141,6 @@ impl Display for Target { pub enum Format { #[clap(name = "csv")] Csv, - #[clap(name = "arrow")] - Arrow, #[clap(name = "parquet")] Parquet, #[clap(name = "vortex")] @@ -189,7 +187,6 @@ impl Format { pub fn name(&self) -> &'static str { match self { Format::Csv => "csv", - Format::Arrow => "arrow", Format::Parquet => "parquet", Format::OnDiskVortex => "vortex-file-compressed", Format::VortexCompact => "vortex-compact", @@ -202,7 +199,6 @@ impl Format { pub fn ext(&self) -> &'static str { match self { Format::Csv => "csv", - Format::Arrow => "arrow", Format::Parquet => "parquet", Format::OnDiskVortex => "vortex", Format::VortexCompact => "vortex", @@ -218,7 +214,6 @@ impl Format { pub enum Engine { #[default] Vortex, - Arrow, #[clap(name = "datafusion")] #[serde(rename = "datafusion")] DataFusion, @@ -233,7 +228,6 @@ impl Display for Engine { Engine::DataFusion => write!(f, "datafusion"), Engine::DuckDB => write!(f, "duckdb"), Engine::Vortex => write!(f, "vortex"), - Engine::Arrow => write!(f, "arrow"), } } } diff --git a/vortex-bench/src/measurements.rs b/vortex-bench/src/measurements.rs index 4f6abe65563..d6d4ad85b32 100644 --- a/vortex-bench/src/measurements.rs +++ b/vortex-bench/src/measurements.rs @@ -352,8 +352,8 @@ impl ToJson for CompressionTimingMeasurement { fn to_json(&self) -> serde_json::Value { let (name, engine) = match self.format { Format::OnDiskVortex => (self.name.to_string(), Engine::Vortex), - Format::Parquet => (format!("parquet_rs-zstd {}", self.name), Engine::Arrow), - Format::Lance => (format!("lance {}", self.name), Engine::Arrow), + Format::Parquet => (format!("parquet_rs-zstd {}", self.name), Engine::Vortex), + Format::Lance => (format!("lance {}", self.name), Engine::Vortex), _ => vortex_panic!( "CompressionTimingMeasurement only supports vortex, lance, and parquet formats" ), @@ -397,8 +397,8 @@ impl ToJson for CustomUnitMeasurement { fn to_json(&self) -> serde_json::Value { let engine = match self.format { Format::OnDiskVortex | Format::VortexCompact => Engine::Vortex, - Format::Parquet => Engine::Arrow, - Format::Lance => Engine::Arrow, + Format::Parquet => Engine::Vortex, + Format::Lance => Engine::Vortex, _ => Engine::Vortex, // Default to Vortex for other formats. }; diff --git a/vortex-bench/src/tpch/tpchgen.rs b/vortex-bench/src/tpch/tpchgen.rs index 79b91200e48..ba2646b365e 100644 --- a/vortex-bench/src/tpch/tpchgen.rs +++ b/vortex-bench/src/tpch/tpchgen.rs @@ -193,7 +193,7 @@ fn generate_table_files( options: TpchGenOptions, ) -> Result>>> { let write_format = match options.format { - Format::Parquet | Format::Arrow | Format::OnDiskDuckDB => Format::Parquet, + Format::Parquet | Format::OnDiskDuckDB => Format::Parquet, Format::OnDiskVortex => Format::OnDiskVortex, Format::VortexCompact => Format::VortexCompact, f => { diff --git a/vortex-bench/src/v3.rs b/vortex-bench/src/v3.rs index 4bdd67ed18f..eb48173224e 100644 --- a/vortex-bench/src/v3.rs +++ b/vortex-bench/src/v3.rs @@ -141,7 +141,7 @@ pub struct QueryMeasurementRecord { pub query_idx: u32, /// Storage backend the run targeted (`nvme` or `s3`). pub storage: String, - /// Query engine (`datafusion`, `duckdb`, `vortex`, `arrow`). + /// Query engine (`datafusion`, `duckdb`, `vortex`). pub engine: String, /// On-disk format (`parquet`, `vortex-file-compressed`, `lance`, ...). pub format: String, @@ -519,7 +519,6 @@ fn duration_as_ns(d: std::time::Duration) -> u64 { fn engine_label(engine: Engine) -> &'static str { match engine { Engine::Vortex => "vortex", - Engine::Arrow => "arrow", Engine::DataFusion => "datafusion", Engine::DuckDB => "duckdb", } @@ -663,7 +662,7 @@ mod tests { fn snapshot_random_access_time() -> anyhow::Result<()> { let timing = TimingMeasurement { name: "random-access/taxi/uniform/parquet-tokio-local-disk".to_string(), - target: Target::new(Engine::Arrow, Format::Parquet), + target: Target::new(Engine::Vortex, Format::Parquet), storage: "nvme".to_string(), runs: vec![ Duration::from_nanos(800_000),