From b826bf157d69e72843c07ea402e06a1d87112b5d Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 7 Oct 2026 10:15:52 +0000 Subject: [PATCH 1/3] Add the JSONBench benchmark suite ClickHouse's JSONBench runs five analytical queries over Bluesky events stored as semi-structured JSON. Add it to vortex-bench as `jsonbench`, with the same events stored three ways: Parquet with the event as a JSON string, Parquet with a shredded Parquet Variant column, and Vortex with a shredded Vortex Variant column converted from the Parquet Variant files. Queries read JSON paths through `{str:a.b}` and `{i64:a.b}` placeholders that each engine expands for each format, through the new `Benchmark::query_for`. DataFusion gets datafusion-functions-json for the JSON strings. Data generation and the DuckDB runner now run a suite's own `prepare_format` before the generic Parquet to Vortex conversion, so JSONBench writes its Vortex files from the Parquet Variant files. The suite is not in the CI benchmark matrix. DataFusion over the Variant formats and DuckDB over Vortex need Variant support in those engines, which lands separately. Signed-off-by: Joe Isaacs Co-Authored-By: Claude Claude-Session: https://claude.ai/code/session_01P9G2RwAZnBNHUackdtcjeq --- Cargo.lock | 37 ++ Cargo.toml | 2 + bench-orchestrator/README.md | 2 +- .../bench_orchestrator/ci_matrix/render.py | 1 + .../bench_orchestrator/config.py | 4 + benchmarks/datafusion-bench/Cargo.toml | 1 + benchmarks/datafusion-bench/src/lib.rs | 16 +- benchmarks/datafusion-bench/src/main.rs | 6 +- benchmarks/duckdb-bench/src/lib.rs | 1 + benchmarks/duckdb-bench/src/main.rs | 6 +- vortex-bench/Cargo.toml | 9 +- vortex-bench/sql/jsonbench.md | 66 ++++ vortex-bench/sql/jsonbench.sql | 49 +++ vortex-bench/src/benchmark.rs | 8 + vortex-bench/src/bin/data-gen.rs | 6 + vortex-bench/src/conversions.rs | 14 + vortex-bench/src/datasets/mod.rs | 5 + vortex-bench/src/jsonbench/data.rs | 323 ++++++++++++++++ vortex-bench/src/jsonbench/mod.rs | 357 ++++++++++++++++++ vortex-bench/src/lib.rs | 17 +- vortex-bench/src/runner.rs | 10 + vortex-bench/src/v3.rs | 15 + 22 files changed, 947 insertions(+), 8 deletions(-) create mode 100644 vortex-bench/sql/jsonbench.md create mode 100644 vortex-bench/sql/jsonbench.sql create mode 100644 vortex-bench/src/jsonbench/data.rs create mode 100644 vortex-bench/src/jsonbench/mod.rs diff --git a/Cargo.lock b/Cargo.lock index 7c3a95148fa..e2128bbcfa8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2799,6 +2799,7 @@ dependencies = [ "custom-labels", "datafusion 55.1.0", "datafusion-common 55.1.0", + "datafusion-functions-json", "datafusion-physical-plan 55.1.0", "futures", "itertools 0.14.0", @@ -3502,6 +3503,19 @@ dependencies = [ "datafusion-physical-expr-common 55.1.0", ] +[[package]] +name = "datafusion-functions-json" +version = "0.55.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f3edcc7d53e84d8a9e05e50b58a153f616ea11d5a07eabcdf5559b532a14119b" +dependencies = [ + "datafusion 55.1.0", + "jiter", + "log", + "paste", + "serde_json", +] + [[package]] name = "datafusion-functions-nested" version = "54.1.0" @@ -5894,6 +5908,21 @@ dependencies = [ "jiff-tzdb", ] +[[package]] +name = "jiter" +version = "0.17.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bc2f2b4e673d798b3dbad0af24bc23b92013e680a6b9ad300b4e5e52fafe0586" +dependencies = [ + "ahash", + "bitvec", + "lexical-parse-float", + "num-bigint 0.4.8", + "num-traits", + "pyo3", + "smallvec", +] + [[package]] name = "jni" version = "0.22.4" @@ -8205,6 +8234,9 @@ dependencies = [ "num-integer", "num-traits", "object_store 0.13.2", + "parquet-variant", + "parquet-variant-compute", + "parquet-variant-json", "seq-macro", "simdutf8", "snap", @@ -8674,6 +8706,8 @@ dependencies = [ "chrono", "indexmap 2.14.2", "libc", + "num-bigint 0.4.8", + "num-traits", "once_cell", "portable-atomic", "pyo3-build-config", @@ -11451,6 +11485,7 @@ dependencies = [ "bytes", "bzip2", "clap", + "flate2", "futures", "geo", "geo-traits", @@ -11469,6 +11504,8 @@ dependencies = [ "parking_lot", "parquet 56.2.1", "parquet 59.3.0", + "parquet-variant", + "parquet-variant-compute", "rand 0.10.3", "regex", "reqwest 0.13.5", diff --git a/Cargo.toml b/Cargo.toml index f26c25fc0f6..509aeabac8d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -148,6 +148,7 @@ datafusion-datasource = { version = "55.0.0", default-features = false } datafusion-execution = { version = "55.0.0" } datafusion-expr = { version = "55.0.0" } datafusion-functions = { version = "55.0.0" } +datafusion-functions-json = "0.55" datafusion-functions-nested = { version = "55.0.0" } datafusion-physical-expr = { version = "55.0.0" } datafusion-physical-expr-adapter = { version = "55.0.0" } @@ -161,6 +162,7 @@ env_logger = "0.11" fastlanes = { version = "0.7.0", features = ["runtime"] } fearless_simd = "1.0.0" fearless_simd_macros = "0.1.0" +flate2 = "1.1" flatbuffers = "25.2.10" fsst-rs = "0.6.0" futures = { version = "0.3.31", default-features = false } diff --git a/bench-orchestrator/README.md b/bench-orchestrator/README.md index 4ecb233955e..1a3704bb4e3 100644 --- a/bench-orchestrator/README.md +++ b/bench-orchestrator/README.md @@ -41,7 +41,7 @@ vx-bench run [options] **Arguments:** -- `benchmark`: Benchmark suite to run (`appian`, `tpch`, `tpcds`, `clickbench`, `fineweb`, `gh-archive`, `polarsignals`, `public-bi`, `statpopgen`) +- `benchmark`: Benchmark suite to run (`appian`, `tpch`, `tpcds`, `clickbench`, `fineweb`, `gh-archive`, `jsonbench`, `polarsignals`, `public-bi`, `statpopgen`) **Options:** diff --git a/bench-orchestrator/bench_orchestrator/ci_matrix/render.py b/bench-orchestrator/bench_orchestrator/ci_matrix/render.py index 8b814349a40..baae03e857b 100644 --- a/bench-orchestrator/bench_orchestrator/ci_matrix/render.py +++ b/bench-orchestrator/bench_orchestrator/ci_matrix/render.py @@ -10,6 +10,7 @@ _NOT_GENERATED = frozenset({Format.LANCE}) _FORMAT_ORDER = ( Format.PARQUET, + Format.PARQUET_VARIANT, Format.VORTEX, Format.VORTEX_COMPACT, Format.VORTEX_SPATIAL_NATIVE, diff --git a/bench-orchestrator/bench_orchestrator/config.py b/bench-orchestrator/bench_orchestrator/config.py index 10cbcfea81d..628d13994f8 100644 --- a/bench-orchestrator/bench_orchestrator/config.py +++ b/bench-orchestrator/bench_orchestrator/config.py @@ -41,6 +41,7 @@ class Format(Enum): """Data formats for benchmarks.""" PARQUET = "parquet" + PARQUET_VARIANT = "parquet-variant" VORTEX = "vortex" VORTEX_COMPACT = "vortex-compact" VORTEX_SPATIAL_NATIVE = "vortex-spatial-native" @@ -58,6 +59,7 @@ class Benchmark(Enum): CLICKBENCH_SORTED = "clickbench-sorted" FINEWEB = "fineweb" GHARCHIVE = "gh-archive" + JSONBENCH = "jsonbench" POLARSIGNALS = "polarsignals" PUBLIC_BI = "public-bi" STATPOPGEN = "statpopgen" @@ -69,12 +71,14 @@ class Benchmark(Enum): ENGINE_FORMATS: dict[Engine, list[Format]] = { Engine.DATAFUSION: [ Format.PARQUET, + Format.PARQUET_VARIANT, Format.VORTEX, Format.VORTEX_COMPACT, Format.LANCE, ], Engine.DUCKDB: [ Format.PARQUET, + Format.PARQUET_VARIANT, Format.VORTEX, Format.VORTEX_COMPACT, Format.VORTEX_SPATIAL_NATIVE, diff --git a/benchmarks/datafusion-bench/Cargo.toml b/benchmarks/datafusion-bench/Cargo.toml index ad0592ffc54..52e6f646e70 100644 --- a/benchmarks/datafusion-bench/Cargo.toml +++ b/benchmarks/datafusion-bench/Cargo.toml @@ -32,6 +32,7 @@ datafusion = { workspace = true, features = [ "unicode_expressions", ] } datafusion-common = { workspace = true } +datafusion-functions-json = { workspace = true } datafusion-physical-plan = { workspace = true } futures.workspace = true itertools.workspace = true diff --git a/benchmarks/datafusion-bench/src/lib.rs b/benchmarks/datafusion-bench/src/lib.rs index 4eec073e7fc..85289e244a4 100644 --- a/benchmarks/datafusion-bench/src/lib.rs +++ b/benchmarks/datafusion-bench/src/lib.rs @@ -21,6 +21,8 @@ use object_store::aws::AmazonS3Builder; use object_store::gcp::GoogleCloudStorageBuilder; use object_store::local::LocalFileSystem; use url::Url; +use vortex_bench::Benchmark; +use vortex_bench::BenchmarkDataset; use vortex_bench::Format; use vortex_bench::SESSION; use vortex_datafusion::VortexFormat; @@ -58,6 +60,18 @@ pub fn get_session_context() -> SessionContext { SessionContext::new_with_state(session_state_builder.build()) } +/// Register the scalar functions `benchmark`'s queries call beyond DataFusion's defaults. +pub fn register_benchmark_functions( + session: &mut SessionContext, + benchmark: &dyn Benchmark, +) -> anyhow::Result<()> { + if matches!(benchmark.dataset(), BenchmarkDataset::JsonBench { .. }) { + // `json_get_*` read the JSON strings of the Parquet baseline. + datafusion_functions_json::register_all(session)?; + } + Ok(()) +} + pub fn make_object_store( session: &SessionContext, source: &Url, @@ -101,7 +115,7 @@ pub fn make_object_store( pub fn format_to_df_format(format: Format) -> anyhow::Result> { Ok(match format { Format::Csv => Arc::new(CsvFormat::default()) as _, - Format::Parquet => Arc::new(ParquetFormat::new()), + Format::Parquet | Format::ParquetVariant => Arc::new(ParquetFormat::new()), Format::OnDiskVortex | Format::VortexCompact | Format::VortexSpatialNative => 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 55f2a03fd08..e93580911f6 100644 --- a/benchmarks/datafusion-bench/src/main.rs +++ b/benchmarks/datafusion-bench/src/main.rs @@ -181,7 +181,8 @@ async fn main() -> anyhow::Result<()> { |format| { let benchmark = &*benchmark; async move { - let session = datafusion_bench::get_session_context(); + let mut session = datafusion_bench::get_session_context(); + datafusion_bench::register_benchmark_functions(&mut session, benchmark)?; for sql in benchmark.engine_init_sql(Engine::DataFusion) { session.sql(&sql).await?.collect().await?; } @@ -192,13 +193,14 @@ async fn main() -> anyhow::Result<()> { }, |query_idx, (session, format), query| { let plans = Arc::clone(&collected_plans); + let query = benchmark.query_for(Engine::DataFusion, *format, query); let labelset = set_labels(benchmark_name.clone(), query_idx, *format); Box::pin( async move { let timer = Instant::now(); - let (batches, plan) = execute_query(session, query) + let (batches, plan) = execute_query(session, &query) .with_labelset(get_labelset_from_global()) .await?; let time = timer.elapsed(); diff --git a/benchmarks/duckdb-bench/src/lib.rs b/benchmarks/duckdb-bench/src/lib.rs index 247dd5351db..f083be9c87b 100644 --- a/benchmarks/duckdb-bench/src/lib.rs +++ b/benchmarks/duckdb-bench/src/lib.rs @@ -208,6 +208,7 @@ impl DuckClient { let object_type = match file_format { Format::Parquet + | Format::ParquetVariant | Format::OnDiskVortex | Format::VortexCompact | Format::VortexSpatialNative => "VIEW", diff --git a/benchmarks/duckdb-bench/src/main.rs b/benchmarks/duckdb-bench/src/main.rs index c680d62f836..0ed300a123b 100644 --- a/benchmarks/duckdb-bench/src/main.rs +++ b/benchmarks/duckdb-bench/src/main.rs @@ -142,6 +142,9 @@ fn main() -> anyhow::Result<()> { .map_err(|_| anyhow::anyhow!("Invalid file URL: {}", benchmark.data_url()))?; for format in args.formats.iter().copied() { + // Benchmark-specific preparation runs first so a suite can write its own Vortex + // files, which the generic Parquet conversion below then leaves alone. + benchmark.prepare_format(format, &base_path).await?; match format { Format::OnDiskVortex => { convert_parquet_directory_to_vortex( @@ -160,7 +163,6 @@ fn main() -> anyhow::Result<()> { // OnDiskDuckDB tables are created during register_tables by loading from Parquet _ => {} } - benchmark.prepare_format(format, &base_path).await?; } anyhow::Ok(()) @@ -222,7 +224,7 @@ fn main() -> anyhow::Result<()> { if !args.reuse { ctx.reopen()?; } - ctx.execute_query_result(query) + ctx.execute_query_result(&benchmark.query_for(Engine::DuckDB, format, query)) }, )?; diff --git a/vortex-bench/Cargo.toml b/vortex-bench/Cargo.toml index 6fac5b74c22..a3d822ad509 100644 --- a/vortex-bench/Cargo.toml +++ b/vortex-bench/Cargo.toml @@ -38,6 +38,7 @@ async-trait = { workspace = true } bytes = { workspace = true } bzip2 = { workspace = true } clap = { workspace = true, features = ["derive"] } +flate2 = { workspace = true } futures = { workspace = true } geo = { workspace = true } geo-traits = { workspace = true } @@ -53,7 +54,13 @@ noodles-bgzf = { workspace = true, features = ["async"] } noodles-vcf = { workspace = true, features = ["async"] } object_store = { workspace = true, features = ["aws"] } parking_lot = { workspace = true } -parquet = { workspace = true, features = ["async", "object_store"] } +parquet = { workspace = true, features = [ + "async", + "object_store", + "variant_experimental", +] } +parquet-variant = { workspace = true } +parquet-variant-compute = { workspace = true } rand = { workspace = true } regex = { workspace = true } reqwest = { workspace = true, features = ["stream"] } diff --git a/vortex-bench/sql/jsonbench.md b/vortex-bench/sql/jsonbench.md new file mode 100644 index 00000000000..a260618e141 --- /dev/null +++ b/vortex-bench/sql/jsonbench.md @@ -0,0 +1,66 @@ +# JSONBench + +ClickHouse's [JSONBench](https://github.com/ClickHouse/JSONBench): analytical queries over +semi-structured JSON, run against Bluesky social network events. Each event is a JSON document +whose shape depends on its type, so this is the main semi-structured / Variant workload. + +The `bluesky` table has a single `data` column holding the whole event. The five queries in +[`jsonbench.sql`](./jsonbench.sql) (numbered from Q0 in file order) are JSONBench's: they read a +handful of JSON paths, filter on some of them and aggregate. The harness lives in +[`src/jsonbench`](../src/jsonbench). + +## Formats + +Every format is derived from the same raw JSON lines, one output file per one-million-event +input file: + +| Format | `data` column | +|---|---| +| `parquet` | the event as a JSON string | +| `parquet-variant` | the event as a shredded Parquet Variant | +| `vortex` | the event as a shredded Vortex Variant, converted from `parquet-variant` | + +Both Variant formats shred the same paths. Data generation infers them from the first 100,000 +events of the first input file, independently of the queries: every scalar path present in at least +1% of the events whose values share one type in at least 99% of its occurrences. The inferred +paths are saved to `shredding.json` next to the data. + +## Queries per engine + +A query reads JSON path `a.b` with `{str:a.b}` (a string) or `{i64:a.b}` (a 64-bit integer). The +harness expands these into each engine's idiom for the format: + +| Engine | `parquet` | `parquet-variant`, `vortex` | +|---|---|---| +| DataFusion | `json_get_str(data, 'a', 'b')` ([datafusion-functions-json]) | `variant_get(data, 'a.b', 'Utf8')` | +| DuckDB | `json_extract_string(data, '$.a.b')` | `CAST(data."a"."b" AS VARCHAR)` | + +DuckDB divides integers into doubles, so Q4's `activity_span_ms` has a fractional part on DuckDB +and is an integer on DataFusion. + +[datafusion-functions-json]: https://github.com/datafusion-contrib/datafusion-functions-json + +## Status + +The suite is not part of the CI benchmark matrix yet. Reading Variant columns needs engine support +that is not in place yet, so these targets fail for now: + +- DataFusion over `parquet-variant` and `vortex`: DataFusion has no `variant_get` function + registered. +- DuckDB over `vortex`: `read_vortex` does not yet expose Vortex Variant columns. + +## Running locally + +The full dataset is 1,000 files of one million events (about 130 GB of compressed JSON). The +`scale-factor` option picks how many millions of events to use (default 1): + +```bash +cargo run --release --bin data-gen -- jsonbench --opt scale-factor=10 \ + --formats parquet,parquet-variant,vortex +cargo run --release --bin datafusion-bench -- jsonbench --opt scale-factor=10 \ + --formats parquet,parquet-variant,vortex +cargo run --release --bin duckdb-bench -- jsonbench --opt scale-factor=10 \ + --formats parquet,parquet-variant,vortex +``` + +Set `VX_BENCH_PRINT_RESULTS=1` to print each query's result, for comparing formats. diff --git a/vortex-bench/sql/jsonbench.sql b/vortex-bench/sql/jsonbench.sql new file mode 100644 index 00000000000..d329c377b27 --- /dev/null +++ b/vortex-bench/sql/jsonbench.sql @@ -0,0 +1,49 @@ +-- JSONBench queries over Bluesky events, numbered from Q0 in file order. The harness splits the +-- file on semicolons, so a comment must never contain one. +-- +-- `{str:a.b}` and `{i64:a.b}` read JSON path `a.b` of the `data` column as a string or a 64-bit +-- integer. The harness expands them into each engine's idiom for each storage format. + +-- Q0. Top event types. +SELECT {str:commit.collection} AS event, count(*) AS cnt +FROM bluesky +GROUP BY event +ORDER BY cnt DESC; + +-- Q1. Top event types together with unique users per event type. +SELECT {str:commit.collection} AS event, count(*) AS cnt, count(DISTINCT {str:did}) AS users +FROM bluesky +WHERE {str:kind} = 'commit' AND {str:commit.operation} = 'create' +GROUP BY event +ORDER BY cnt DESC; + +-- Q2. When do people use Bluesky. +SELECT {str:commit.collection} AS event, + date_part('hour', to_timestamp({i64:time_us} / 1000000)) AS hour_of_day, + count(*) AS cnt +FROM bluesky +WHERE {str:kind} = 'commit' + AND {str:commit.operation} = 'create' + AND {str:commit.collection} IN ('app.bsky.feed.post', 'app.bsky.feed.repost', 'app.bsky.feed.like') +GROUP BY event, hour_of_day +ORDER BY hour_of_day, event; + +-- Q3. The three users who posted first. +SELECT {str:did} AS user_id, min({i64:time_us}) AS first_post_us +FROM bluesky +WHERE {str:kind} = 'commit' + AND {str:commit.operation} = 'create' + AND {str:commit.collection} = 'app.bsky.feed.post' +GROUP BY user_id +ORDER BY first_post_us ASC, user_id +LIMIT 3; + +-- Q4. The three users with the longest posting activity span. +SELECT {str:did} AS user_id, (max({i64:time_us}) - min({i64:time_us})) / 1000 AS activity_span_ms +FROM bluesky +WHERE {str:kind} = 'commit' + AND {str:commit.operation} = 'create' + AND {str:commit.collection} = 'app.bsky.feed.post' +GROUP BY user_id +ORDER BY activity_span_ms DESC, user_id +LIMIT 3; diff --git a/vortex-bench/src/benchmark.rs b/vortex-bench/src/benchmark.rs index 47de30a5faf..86e73653f4a 100644 --- a/vortex-bench/src/benchmark.rs +++ b/vortex-bench/src/benchmark.rs @@ -41,6 +41,14 @@ pub trait Benchmark: Send + Sync { Vec::new() } + /// The SQL `engine` runs for `query` against `format`. + /// + /// Suites whose SQL depends on how a format stores the data (e.g. JSON strings versus Variant + /// columns) rewrite it here. Default: `query` unchanged. + fn query_for(&self, _engine: Engine, _format: Format, query: &str) -> String { + query.to_owned() + } + /// Generate or prepare base data for the benchmark (typically Parquet format). /// This is the canonical source data that can be converted to other formats. /// This should be idempotent - safe to call multiple times. diff --git a/vortex-bench/src/bin/data-gen.rs b/vortex-bench/src/bin/data-gen.rs index 35c77d70c48..ccb4ae221b6 100644 --- a/vortex-bench/src/bin/data-gen.rs +++ b/vortex-bench/src/bin/data-gen.rs @@ -69,6 +69,12 @@ async fn main() -> anyhow::Result<()> { .to_file_path() .map_err(|_| anyhow::anyhow!("Invalid file URL: {}", benchmark.data_url()))?; + // Benchmark-specific preparation runs first so a suite can write its own Vortex files, + // which the generic Parquet conversion below then leaves alone. + for format in args.formats.iter().copied() { + benchmark.prepare_format(format, &base_path).await?; + } + if args .formats .iter() diff --git a/vortex-bench/src/conversions.rs b/vortex-bench/src/conversions.rs index fda1638a63e..4b8f80e7c18 100644 --- a/vortex-bench/src/conversions.rs +++ b/vortex-bench/src/conversions.rs @@ -154,6 +154,10 @@ pub async fn parquet_to_vortex_chunks_with_batch_size( fn record_batch_to_vortex(batch: RecordBatch) -> VortexResult { let schema = batch.schema(); let chunk = SESSION.arrow().from_arrow_record_batch(batch, &schema)?; + // Variant arrays have no builder. The writer's compressor canonicalizes them instead. + if contains_variant(chunk.dtype()) { + return Ok(chunk); + } let mut ctx = VortexSession::default().create_execution_ctx(); let mut builder = builder_with_capacity_in(chunk.dtype(), chunk.len(), ctx.allocator()); @@ -163,6 +167,16 @@ fn record_batch_to_vortex(batch: RecordBatch) -> VortexResult { Ok(builder.finish()) } +/// Whether `dtype` is or nests a Variant. +fn contains_variant(dtype: &DType) -> bool { + match dtype { + DType::Variant(_) => true, + DType::Struct(fields, _) => fields.fields().any(|field| contains_variant(&field)), + DType::List(element, _) | DType::FixedSizeList(element, ..) => contains_variant(element), + _ => false, + } +} + /// Create a streaming Vortex array from a Parquet reader. /// /// Streams record batches and converts them to Vortex arrays on-the-fly, avoiding loading the diff --git a/vortex-bench/src/datasets/mod.rs b/vortex-bench/src/datasets/mod.rs index 8133af048f8..ff2d4b6036c 100644 --- a/vortex-bench/src/datasets/mod.rs +++ b/vortex-bench/src/datasets/mod.rs @@ -82,6 +82,8 @@ pub enum BenchmarkDataset { Fineweb, #[serde(rename = "gharchive")] GhArchive, + #[serde(rename = "jsonbench")] + JsonBench { n_rows: usize }, #[serde(rename = "vortex")] VortexQueries, } @@ -100,6 +102,7 @@ impl BenchmarkDataset { BenchmarkDataset::PolarSignals { .. } => "polarsignals", BenchmarkDataset::Fineweb => "fineweb", BenchmarkDataset::GhArchive => "gharchive", + BenchmarkDataset::JsonBench { .. } => "jsonbench", BenchmarkDataset::VortexQueries => "vortex", } } @@ -126,6 +129,7 @@ impl Display for BenchmarkDataset { } BenchmarkDataset::Fineweb => write!(f, "fineweb"), BenchmarkDataset::GhArchive => write!(f, "gharchive"), + BenchmarkDataset::JsonBench { n_rows } => write!(f, "jsonbench(n_rows={n_rows})"), BenchmarkDataset::VortexQueries => write!(f, "vortex"), } } @@ -186,6 +190,7 @@ impl BenchmarkDataset { BenchmarkDataset::PolarSignals { .. } => &["stacktraces"], BenchmarkDataset::Fineweb => &["fineweb"], BenchmarkDataset::GhArchive => &["events"], + BenchmarkDataset::JsonBench { .. } => &["bluesky"], // See VortexBenchmark::table_specs BenchmarkDataset::VortexQueries => &[], } diff --git a/vortex-bench/src/jsonbench/data.rs b/vortex-bench/src/jsonbench/data.rs new file mode 100644 index 00000000000..2c0dca9de91 --- /dev/null +++ b/vortex-bench/src/jsonbench/data.rs @@ -0,0 +1,323 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! Data preparation for JSONBench. +//! +//! The raw data is ClickHouse's JSONBench Bluesky event dump: gzipped JSON-lines files of one +//! million events each. Every format is derived from the same raw lines: +//! +//! - `parquet`: one `data` column holding each event as a JSON string. +//! - `parquet-variant`: one `data` column holding each event as a shredded Parquet Variant. +//! - `vortex`/`vortex-compact`: the `parquet-variant` files converted to Vortex, so the Vortex +//! `Variant` column carries exactly the same shredded values. + +use std::collections::BTreeMap; +use std::fs::File; +use std::io::BufRead; +use std::io::BufReader; +use std::path::Path; +use std::sync::Arc; + +use anyhow::Context; +use arrow_array::ArrayRef; +use arrow_array::RecordBatch; +use arrow_array::StringArray; +use arrow_schema::DataType; +use arrow_schema::Field; +use arrow_schema::Schema; +use arrow_schema::SchemaRef; +use flate2::read::MultiGzDecoder; +use parquet::arrow::ArrowWriter; +use parquet::basic::Compression; +use parquet::basic::ZstdLevel; +use parquet::file::properties::WriterProperties; +use parquet_variant::VariantPath; +use parquet_variant::VariantPathElement; +use parquet_variant_compute::ShreddedSchemaBuilder; +use parquet_variant_compute::json_to_variant; +use parquet_variant_compute::shred_variant; +use serde::Deserialize; +use serde::Serialize; +use serde_json::Value; +use tracing::info; +use tracing::warn; + +/// Rows in each raw JSONBench file. +pub const ROWS_PER_FILE: usize = 1_000_000; + +/// Rows per record batch while converting. +const BATCH_ROWS: usize = 65_536; + +/// Rows sampled from the first raw file to infer the shredding schema. +const SHREDDING_SAMPLE_ROWS: usize = 100_000; + +/// Minimum fraction of sampled rows a path must appear in to be shredded. +const SHREDDING_MIN_PRESENCE: f64 = 0.01; + +/// Minimum fraction of a path's occurrences that must share one scalar type for it to be shredded. +const SHREDDING_MIN_TYPE_SHARE: f64 = 0.99; + +/// Deepest object nesting the shredding inference descends into. +const SHREDDING_MAX_DEPTH: usize = 6; + +/// URL of the `file_idx`-th (1-based) raw JSONBench file. +pub fn raw_json_url(file_idx: usize) -> String { + format!( + "https://clickhouse-public-datasets.s3.amazonaws.com/bluesky/file_{file_idx:04}.json.gz" + ) +} + +/// File name of the `file_idx`-th (1-based) raw JSONBench file. +pub fn raw_json_name(file_idx: usize) -> String { + format!("file_{file_idx:04}.json.gz") +} + +/// Stem shared by every derived file of the `file_idx`-th raw file. +pub fn output_stem(file_idx: usize) -> String { + format!("bluesky_{file_idx:04}") +} + +/// The scalar type a shredded path is stored as. +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum ShredType { + Boolean, + Int64, + Float64, + Utf8, +} + +impl ShredType { + fn of(value: &Value) -> Option { + match value { + Value::Bool(_) => Some(Self::Boolean), + Value::Number(n) if n.is_i64() => Some(Self::Int64), + Value::Number(_) => Some(Self::Float64), + Value::String(_) => Some(Self::Utf8), + Value::Null | Value::Array(_) | Value::Object(_) => None, + } + } + + fn data_type(self) -> DataType { + match self { + Self::Boolean => DataType::Boolean, + Self::Int64 => DataType::Int64, + Self::Float64 => DataType::Float64, + Self::Utf8 => DataType::Utf8, + } + } +} + +/// One shredded path: object field names from the root, and the type the leaf is stored as. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct ShreddedPath { + pub path: Vec, + pub shred_type: ShredType, +} + +/// The paths shredded into typed columns by both Variant formats. +#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct ShreddingSchema { + pub paths: Vec, +} + +impl ShreddingSchema { + /// Infer the shredding schema from the first `SHREDDING_SAMPLE_ROWS` rows of `raw`. + /// + /// A scalar leaf is shredded when it appears in at least `SHREDDING_MIN_PRESENCE` of the sampled + /// rows and at least `SHREDDING_MIN_TYPE_SHARE` of its occurrences share one type. This mirrors + /// the automatic path discovery engines with native JSON storage perform, and is independent of + /// the benchmark queries. + pub fn infer(raw: &Path) -> anyhow::Result { + let mut stats: BTreeMap, BTreeMap, usize>> = BTreeMap::new(); + let mut rows = 0usize; + for line in json_lines(raw)? { + let line = line?; + let Ok(value) = serde_json::from_str::(&line) else { + continue; + }; + rows += 1; + let mut path = Vec::new(); + collect_leaf_types(&value, &mut path, &mut stats); + if rows == SHREDDING_SAMPLE_ROWS { + break; + } + } + anyhow::ensure!(rows > 0, "no JSON rows in {}", raw.display()); + + let mut paths = Vec::new(); + for (path, types) in stats { + let occurrences: usize = types.values().sum(); + let Some((Some(shred_type), &count)) = types.iter().max_by_key(|(_, count)| **count) + else { + continue; + }; + #[expect(clippy::cast_precision_loss)] + let presence = occurrences as f64 / rows as f64; + #[expect(clippy::cast_precision_loss)] + let type_share = count as f64 / occurrences as f64; + if presence >= SHREDDING_MIN_PRESENCE && type_share >= SHREDDING_MIN_TYPE_SHARE { + paths.push(ShreddedPath { + path, + shred_type: *shred_type, + }); + } + } + Ok(Self { paths }) + } + + /// The Arrow shredding type passed to [`shred_variant`]. + pub fn arrow_shredding_type(&self) -> anyhow::Result { + let mut builder = ShreddedSchemaBuilder::new(); + for shredded in &self.paths { + let path = VariantPath::from_iter( + shredded + .path + .iter() + .map(|name| VariantPathElement::from(name.as_str())), + ); + builder = builder.with_path(path, &shredded.shred_type.data_type())?; + } + Ok(builder.build()) + } + + pub fn load(path: &Path) -> anyhow::Result { + Ok(serde_json::from_reader(BufReader::new(File::open(path)?))?) + } + + pub fn save(&self, path: &Path) -> anyhow::Result<()> { + serde_json::to_writer_pretty(File::create(path)?, self)?; + Ok(()) + } +} + +/// Record the type of every scalar leaf under `value`, keyed by its object path. +/// +/// Values inside arrays are not descended into: Parquet Variant shredding cannot address list +/// elements by path. +fn collect_leaf_types( + value: &Value, + path: &mut Vec, + stats: &mut BTreeMap, BTreeMap, usize>>, +) { + if let Value::Object(fields) = value + && path.len() < SHREDDING_MAX_DEPTH + { + // An object is not a scalar leaf, but counting it keeps a path that holds an object in + // some rows and a scalar in others from being shredded as that scalar. + if !path.is_empty() { + *stats + .entry(path.clone()) + .or_default() + .entry(None) + .or_default() += 1; + } + for (name, child) in fields { + path.push(name.clone()); + collect_leaf_types(child, path, stats); + path.pop(); + } + } else if !path.is_empty() { + *stats + .entry(path.clone()) + .or_default() + .entry(ShredType::of(value)) + .or_default() += 1; + } +} + +/// Iterate the lines of a gzipped JSON-lines file. +fn json_lines(raw: &Path) -> anyhow::Result>> { + let file = File::open(raw).with_context(|| format!("opening {}", raw.display()))?; + Ok(BufReader::with_capacity(1 << 20, MultiGzDecoder::new(file)).lines()) +} + +/// Batches of the valid JSON documents in `raw`, as a `data` string column. +/// +/// Lines that do not parse as JSON are skipped, so every format holds the same rows. +fn json_batches(raw: &Path) -> anyhow::Result>> { + let mut lines = json_lines(raw)?; + let mut skipped = 0usize; + let raw_display = raw.display().to_string(); + Ok(std::iter::from_fn(move || { + let mut batch = Vec::with_capacity(BATCH_ROWS); + for line in lines.by_ref() { + let line = match line { + Ok(line) => line, + Err(err) => return Some(Err(err.into())), + }; + if serde_json::from_str::(&line).is_err() { + skipped += 1; + continue; + } + batch.push(line); + if batch.len() == BATCH_ROWS { + break; + } + } + if batch.is_empty() { + if skipped > 0 { + warn!("skipped {skipped} invalid JSON lines in {raw_display}"); + } + return None; + } + Some(Ok(StringArray::from(batch))) + })) +} + +/// Parquet writer properties for every Parquet file this benchmark writes. +fn parquet_writer_properties() -> anyhow::Result { + Ok(WriterProperties::builder() + .set_compression(Compression::ZSTD(ZstdLevel::try_new(3)?)) + .build()) +} + +/// Write the JSON documents of `raw` to `output` as a Parquet file with a `data` string column. +pub fn write_json_parquet(raw: &Path, output: &Path) -> anyhow::Result<()> { + let schema: SchemaRef = Arc::new(Schema::new(vec![Field::new("data", DataType::Utf8, false)])); + let mut writer = ArrowWriter::try_new( + File::create(output)?, + Arc::clone(&schema), + Some(parquet_writer_properties()?), + )?; + for batch in json_batches(raw)? { + let column: ArrayRef = Arc::new(batch?); + writer.write(&RecordBatch::try_new(Arc::clone(&schema), vec![column])?)?; + } + writer.close()?; + info!("wrote {}", output.display()); + Ok(()) +} + +/// Write the JSON documents of `raw` to `output` as a Parquet file with a shredded Variant `data` +/// column. +pub fn write_variant_parquet( + raw: &Path, + output: &Path, + shredding: &ShreddingSchema, +) -> anyhow::Result<()> { + let shredding_type = shredding.arrow_shredding_type()?; + let mut writer: Option> = None; + for batch in json_batches(raw)? { + let strings: ArrayRef = Arc::new(batch?); + let variant = shred_variant(&json_to_variant(&strings)?, &shredding_type)?; + let schema = Arc::new(Schema::new(vec![ + variant.field("data").with_nullable(false), + ])); + let batch = RecordBatch::try_new(Arc::clone(&schema), vec![ArrayRef::from(variant)])?; + let writer = match writer.as_mut() { + Some(writer) => writer, + None => writer.insert(ArrowWriter::try_new( + File::create(output)?, + schema, + Some(parquet_writer_properties()?), + )?), + }; + writer.write(&batch)?; + } + writer + .ok_or_else(|| anyhow::anyhow!("no JSON rows in {}", raw.display()))? + .close()?; + info!("wrote {}", output.display()); + Ok(()) +} diff --git a/vortex-bench/src/jsonbench/mod.rs b/vortex-bench/src/jsonbench/mod.rs new file mode 100644 index 00000000000..eb432d1a8b4 --- /dev/null +++ b/vortex-bench/src/jsonbench/mod.rs @@ -0,0 +1,357 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! [JSONBench](https://github.com/ClickHouse/JSONBench): ClickHouse's analytical benchmark over +//! semi-structured JSON, run against Bluesky social network events. +//! +//! Every event lives in a single `data` column of the `bluesky` table. The baseline Parquet format +//! stores it as a JSON string; `parquet-variant` and the Vortex formats store it as a shredded +//! Variant. The queries in `sql/jsonbench.sql` address JSON paths with placeholders that +//! [`JsonBenchBenchmark::query_for`] expands into each engine's idiom for each format. + +pub mod data; + +use std::fs; +use std::path::Path; +use std::path::PathBuf; +use std::sync::LazyLock; +use std::sync::atomic::AtomicUsize; +use std::sync::atomic::Ordering; + +use anyhow::Context; +use regex::Captures; +use regex::Regex; +use tokio::io::AsyncWriteExt; +use tracing::info; +use url::Url; +use vortex::utils::parallelism::get_available_parallelism; + +use crate::Benchmark; +use crate::BenchmarkDataset; +use crate::CompactionStrategy; +use crate::Engine; +use crate::Format; +use crate::TableSpec; +use crate::conversions::convert_parquet_file_to_vortex; +use crate::idempotent; +use crate::idempotent_async; +use crate::jsonbench::data::ShreddingSchema; +use crate::jsonbench::data::output_stem; +use crate::jsonbench::data::raw_json_name; +use crate::jsonbench::data::raw_json_url; +use crate::jsonbench::data::write_json_parquet; +use crate::jsonbench::data::write_variant_parquet; +use crate::utils::file::data_dir; +use crate::workspace_root; + +/// The single table every query reads. +const TABLE: &str = "bluesky"; + +/// Matches a JSON path placeholder such as `{str:commit.collection}` or `{i64:time_us}`. +static PLACEHOLDER: LazyLock = LazyLock::new(|| { + Regex::new(r"\{(str|i64):([A-Za-z0-9_$.]+)\}").expect("placeholder regex is valid") +}); + +pub struct JsonBenchBenchmark { + /// Number of one-million-row raw files the dataset is built from. + n_files: usize, + data_url: Url, + /// Directory holding the downloaded raw files, shared by every scale. + raw_dir: PathBuf, +} + +impl JsonBenchBenchmark { + /// A JSONBench dataset of `scale_factor` million events. + pub fn new(scale_factor: usize) -> anyhow::Result { + anyhow::ensure!(scale_factor > 0, "jsonbench scale factor must be positive"); + let root = data_dir().join("jsonbench"); + let data_dir = root.join(format!("{scale_factor}m")); + let data_url = Url::from_directory_path(&data_dir) + .map_err(|_| anyhow::anyhow!("invalid data directory {}", data_dir.display()))?; + Ok(Self { + n_files: scale_factor, + data_url, + raw_dir: root.join("json"), + }) + } + + fn base_path(&self) -> anyhow::Result { + self.data_url + .to_file_path() + .map_err(|_| anyhow::anyhow!("jsonbench data URL must be a file URL")) + } + + fn raw_path(&self, file_idx: usize) -> PathBuf { + self.raw_dir.join(raw_json_name(file_idx)) + } + + fn file_indices(&self) -> std::ops::RangeInclusive { + 1..=self.n_files + } + + fn shredding_schema(&self) -> anyhow::Result { + let path = self.base_path()?.join("shredding.json"); + idempotent(&path, |tmp| { + let schema = ShreddingSchema::infer(&self.raw_path(1))?; + info!( + "inferred {} shredded JSONBench paths: {}", + schema.paths.len(), + schema + .paths + .iter() + .map(|p| p.path.join(".")) + .collect::>() + .join(", ") + ); + schema.save(tmp) + })?; + ShreddingSchema::load(&path) + } + + /// Run `convert(raw, output)` for every raw file whose `{stem}.{ext}` output is missing in + /// `format_dir`, in parallel. + fn convert_files( + &self, + format_dir: &Path, + ext: &str, + convert: impl Fn(&Path, &Path) -> anyhow::Result<()> + Sync, + ) -> anyhow::Result<()> { + fs::create_dir_all(format_dir)?; + let indices: Vec = self.file_indices().collect(); + let next = AtomicUsize::new(0); + let workers = get_available_parallelism() + .unwrap_or(1) + .min(indices.len()) + .max(1); + std::thread::scope(|scope| { + let handles: Vec<_> = (0..workers) + .map(|_| { + scope.spawn(|| -> anyhow::Result<()> { + while let Some(&file_idx) = + indices.get(next.fetch_add(1, Ordering::Relaxed)) + { + let output = + format_dir.join(format!("{}.{ext}", output_stem(file_idx))); + idempotent(&output, |tmp| convert(&self.raw_path(file_idx), tmp))?; + } + Ok(()) + }) + }) + .collect(); + handles + .into_iter() + .try_for_each(|handle| handle.join().expect("conversion worker panicked")) + }) + } + + fn prepare_variant_parquet(&self) -> anyhow::Result { + let shredding = self.shredding_schema()?; + let dir = self.base_path()?.join(Format::ParquetVariant.name()); + self.convert_files(&dir, "parquet", |raw, output| { + write_variant_parquet(raw, output, &shredding) + })?; + Ok(dir) + } + + async fn prepare_vortex(&self, format: Format) -> anyhow::Result<()> { + let compaction = match format { + Format::VortexCompact => CompactionStrategy::Compact, + _ => CompactionStrategy::Default, + }; + let variant_dir = self.prepare_variant_parquet()?; + let vortex_dir = self.base_path()?.join(format.name()); + fs::create_dir_all(&vortex_dir)?; + for file_idx in self.file_indices() { + let stem = output_stem(file_idx); + let parquet = variant_dir.join(format!("{stem}.parquet")); + idempotent_async(vortex_dir.join(format!("{stem}.vortex")), |tmp| async move { + convert_parquet_file_to_vortex(&parquet, &tmp, compaction).await + }) + .await?; + } + Ok(()) + } + + /// Expand `{str:path}` and `{i64:path}` placeholders into `engine`'s idiom for reading that + /// JSON path from `format`. + pub fn expand_placeholders(engine: Engine, format: Format, query: &str) -> String { + PLACEHOLDER + .replace_all(query, |caps: &Captures| { + let ty = &caps[1]; + let path = &caps[2]; + expand_path(engine, format, ty, path) + }) + .into_owned() + } +} + +/// The SQL expression reading JSON `path` as `ty` (`str` or `i64`) on `engine` from `format`. +fn expand_path(engine: Engine, format: Format, ty: &str, path: &str) -> String { + let string_storage = matches!(format, Format::Parquet | Format::OnDiskDuckDB); + match engine { + Engine::DataFusion if string_storage => { + let keys = path + .split('.') + .map(|key| format!("'{key}'")) + .collect::>() + .join(", "); + match ty { + "str" => format!("json_get_str(data, {keys})"), + _ => format!("json_get_int(data, {keys})"), + } + } + Engine::DataFusion | Engine::Vortex => { + let arrow_type = if ty == "str" { "Utf8" } else { "Int64" }; + format!("variant_get(data, '{path}', '{arrow_type}')") + } + Engine::DuckDB if string_storage => match ty { + "str" => format!("json_extract_string(data, '$.{path}')"), + _ => format!("CAST(json_extract_string(data, '$.{path}') AS BIGINT)"), + }, + Engine::DuckDB => { + let sql_type = if ty == "str" { "VARCHAR" } else { "BIGINT" }; + let quoted = path + .split('.') + .map(|key| format!("\"{key}\"")) + .collect::>() + .join("."); + format!("CAST(data.{quoted} AS {sql_type})") + } + } +} + +#[async_trait::async_trait] +impl Benchmark for JsonBenchBenchmark { + fn queries(&self) -> anyhow::Result> { + // `;`-separated; a `;` must not appear in a comment, or it would split a statement in two. + let queries_file = workspace_root() + .join("vortex-bench") + .join("sql") + .join("jsonbench.sql"); + let contents = fs::read_to_string(&queries_file) + .with_context(|| format!("reading {}", queries_file.display()))?; + Ok(contents + .split_terminator(';') + .map(str::trim) + .filter(|stmt| !stmt.is_empty()) + .map(str::to_string) + .enumerate() + .collect()) + } + + fn query_for(&self, engine: Engine, format: Format, query: &str) -> String { + Self::expand_placeholders(engine, format, query) + } + + async fn generate_base_data(&self) -> anyhow::Result<()> { + fs::create_dir_all(&self.raw_dir)?; + let client = reqwest::Client::new(); + for file_idx in self.file_indices() { + let client = &client; + idempotent_async(self.raw_path(file_idx), |tmp| async move { + let url = raw_json_url(file_idx); + info!("downloading {url}"); + let body = client + .get(&url) + .send() + .await? + .error_for_status() + .with_context(|| format!("fetching {url}"))? + .bytes() + .await?; + let mut file = tokio::fs::File::create(&tmp).await?; + file.write_all(&body).await?; + file.flush().await?; + anyhow::Ok(()) + }) + .await?; + } + + let parquet_dir = self.base_path()?.join(Format::Parquet.name()); + tokio::task::block_in_place(|| { + self.shredding_schema()?; + self.convert_files(&parquet_dir, "parquet", write_json_parquet) + }) + } + + async fn prepare_format(&self, format: Format, _base_path: &Path) -> anyhow::Result<()> { + match format { + Format::ParquetVariant => { + tokio::task::block_in_place(|| self.prepare_variant_parquet())?; + } + Format::OnDiskVortex | Format::VortexCompact => { + tokio::task::block_in_place(|| self.prepare_variant_parquet())?; + self.prepare_vortex(format).await?; + } + _ => {} + } + Ok(()) + } + + fn dataset(&self) -> BenchmarkDataset { + BenchmarkDataset::JsonBench { + n_rows: self.n_files * data::ROWS_PER_FILE, + } + } + + fn doc_path(&self) -> &'static str { + "vortex-bench/sql/jsonbench.md" + } + + fn dataset_name(&self) -> &str { + "jsonbench" + } + + fn dataset_display(&self) -> String { + format!("jsonbench({}m)", self.n_files) + } + + fn data_url(&self) -> &Url { + &self.data_url + } + + fn table_specs(&self) -> Vec { + vec![TableSpec::new(TABLE, None)] + } +} + +#[cfg(test)] +mod tests { + use rstest::rstest; + + use super::*; + + const QUERY: &str = "SELECT {str:commit.collection}, {i64:time_us} FROM bluesky"; + + #[rstest] + #[case( + Engine::DataFusion, + Format::Parquet, + "SELECT json_get_str(data, 'commit', 'collection'), json_get_int(data, 'time_us') FROM bluesky" + )] + #[case( + Engine::DataFusion, + Format::OnDiskVortex, + "SELECT variant_get(data, 'commit.collection', 'Utf8'), variant_get(data, 'time_us', 'Int64') FROM bluesky" + )] + #[case( + Engine::DuckDB, + Format::Parquet, + "SELECT json_extract_string(data, '$.commit.collection'), CAST(json_extract_string(data, '$.time_us') AS BIGINT) FROM bluesky" + )] + #[case( + Engine::DuckDB, + Format::ParquetVariant, + "SELECT CAST(data.\"commit\".\"collection\" AS VARCHAR), CAST(data.\"time_us\" AS BIGINT) FROM bluesky" + )] + fn expands_placeholders( + #[case] engine: Engine, + #[case] format: Format, + #[case] expected: &str, + ) { + assert_eq!( + JsonBenchBenchmark::expand_placeholders(engine, format, QUERY), + expected + ); + } +} diff --git a/vortex-bench/src/lib.rs b/vortex-bench/src/lib.rs index 106a6272411..bbe1ddbe53f 100644 --- a/vortex-bench/src/lib.rs +++ b/vortex-bench/src/lib.rs @@ -17,6 +17,7 @@ use clickbench::ClickBenchSortedBenchmark; use clickbench::Flavor; use fineweb::FinewebBenchmark; use itertools::Itertools; +use jsonbench::JsonBenchBenchmark; use polarsignals::PolarSignalsBenchmark; use public_bi::PBIDataset; use public_bi::PublicBiBenchmark; @@ -47,6 +48,7 @@ pub mod datasets; pub mod display; pub mod downloadable_dataset; pub mod fineweb; +pub mod jsonbench; pub mod measurements; pub mod memory; pub mod output; @@ -145,6 +147,11 @@ pub enum Format { Csv, #[clap(name = "parquet")] Parquet, + /// Parquet with semi-structured columns stored as shredded Parquet Variant rather than JSON + /// strings. Only suites that write it (JSONBench) support it. + #[clap(name = "parquet-variant")] + #[serde(rename = "parquet-variant")] + ParquetVariant, #[clap(name = "vortex")] #[serde(rename = "vortex")] OnDiskVortex, @@ -196,6 +203,7 @@ impl Format { Format::ArrowIpc => "arrow-ipc", Format::Csv => "csv", Format::Parquet => "parquet", + Format::ParquetVariant => "parquet-variant", Format::OnDiskVortex => "vortex-file-compressed", Format::VortexCompact => "vortex-compact", Format::VortexSpatialNative => "vortex-spatial-native", @@ -208,7 +216,7 @@ impl Format { match self { Format::ArrowIpc => "arrow", Format::Csv => "csv", - Format::Parquet => "parquet", + Format::Parquet | Format::ParquetVariant => "parquet", Format::OnDiskVortex => "vortex", Format::VortexCompact => "vortex", Format::VortexSpatialNative => "vortex", @@ -342,6 +350,8 @@ pub enum BenchmarkArg { Fineweb, #[clap(name = "gharchive")] GhArchive, + #[clap(name = "jsonbench")] + JsonBench, #[clap(name = "polarsignals")] PolarSignals, #[clap(name = "public-bi")] @@ -404,6 +414,11 @@ pub fn create_benchmark(b: BenchmarkArg, opts: &Opts) -> anyhow::Result { + let scale_factor = opts.get_as::(SCALE_FACTOR_KEY).unwrap_or(1); + let benchmark = JsonBenchBenchmark::new(scale_factor)?; + Ok(Box::new(benchmark) as _) + } BenchmarkArg::PolarSignals => { let scale_factor = opts.get_as::(SCALE_FACTOR_KEY).unwrap_or(1); let benchmark = PolarSignalsBenchmark::new(scale_factor)?; diff --git a/vortex-bench/src/runner.rs b/vortex-bench/src/runner.rs index 9e5c4db3d08..66ce1bfa08c 100644 --- a/vortex-bench/src/runner.rs +++ b/vortex-bench/src/runner.rs @@ -150,6 +150,7 @@ impl SqlBenchmarkRunner { if row_count.is_none() { row_count = Some(result.row_count()); + print_result_if_requested(query_idx, format, result); } } @@ -411,6 +412,7 @@ impl SqlBenchmarkRunner { if row_count.is_none() { row_count = Some(result.row_count()); + print_result_if_requested(query_idx, format, result); } } @@ -444,6 +446,14 @@ impl SqlBenchmarkRunner { } } +/// Print a query's result to stderr when `VX_BENCH_PRINT_RESULTS=1`, so results can be compared +/// across formats and engines. +fn print_result_if_requested(query_idx: usize, format: Format, result: R) { + if std::env::var("VX_BENCH_PRINT_RESULTS").is_ok_and(|v| v == "1") { + eprintln!("=== Q{query_idx} [{format}] ===\n{}", result.display()); + } +} + fn is_ci() -> bool { matches!(std::env::var("CI").as_deref(), Ok("true")) } diff --git a/vortex-bench/src/v3.rs b/vortex-bench/src/v3.rs index 51a2d71262b..67fb98995ad 100644 --- a/vortex-bench/src/v3.rs +++ b/vortex-bench/src/v3.rs @@ -282,6 +282,7 @@ fn canonical_tpc_scale_factor(scale_factor: &str) -> String { /// | `PolarSignals { n_rows: _ }`| `polarsignals` | `None` | `None` | Same as StatPopGen. | /// | `Fineweb` | `fineweb` | `None` | `None` | | /// | `GhArchive` | `gharchive` | `None` | `None` | | +/// | `JsonBench { n_rows }` | `jsonbench` | `None` | millions of rows as string (`"1"`, `"10"`, ...) | New live-only suite. | /// | `Appian` | `appian` | `None` | `None` | Static dataset; no scale factor. | /// | `PublicBi { name }` | `public-bi` | dataset name (e.g. `cms-provider`) | `None` | Sub-dataset name lives in `dataset_variant`. | /// | `SpatialBench { scale_factor }` | `spatialbench` | `None` | SF as string | Same canonicalization as TPC-H; no historical v2 records to merge with. | @@ -320,6 +321,11 @@ pub fn benchmark_dataset_dims(d: &BenchmarkDataset) -> (String, Option, BenchmarkDataset::PolarSignals { .. } => ("polarsignals".to_string(), None, None), BenchmarkDataset::Fineweb => ("fineweb".to_string(), None, None), BenchmarkDataset::GhArchive => ("gharchive".to_string(), None, None), + BenchmarkDataset::JsonBench { n_rows } => ( + "jsonbench".to_string(), + None, + Some((n_rows / 1_000_000).to_string()), + ), BenchmarkDataset::Appian => ("appian".to_string(), None, None), BenchmarkDataset::VortexQueries => ("vortex".to_string(), None, None), } @@ -744,6 +750,15 @@ mod tests { } } + #[test] + fn jsonbench_dims_carry_millions_of_rows() { + let dims = benchmark_dataset_dims(&BenchmarkDataset::JsonBench { n_rows: 10_000_000 }); + assert_eq!( + dims, + ("jsonbench".to_string(), None, Some("10".to_string())) + ); + } + #[test] fn clickbench_sorted_dims_are_distinct_from_clickbench() { let (dataset, variant, scale_factor) = From b20730ae7cf2208223ff0d42666d6b0faedeaf20 Mon Sep 17 00:00:00 2001 From: Joe Isaacs Date: Wed, 7 Oct 2026 11:00:48 +0000 Subject: [PATCH 2/3] Clean up the JSONBench suite Store the data directory directly instead of re-deriving it from the URL, run Variant Parquet preparation once and off the async runtime when preparing Vortex, make placeholder expansion a free function, and build the raw file URL from its file name. Signed-off-by: Joe Isaacs Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01HXvxbPjYxoaTEqBgPLXbad --- vortex-bench/src/jsonbench/data.rs | 3 +- vortex-bench/src/jsonbench/mod.rs | 52 ++++++++++++------------------ 2 files changed, 22 insertions(+), 33 deletions(-) diff --git a/vortex-bench/src/jsonbench/data.rs b/vortex-bench/src/jsonbench/data.rs index 2c0dca9de91..4c0d8459678 100644 --- a/vortex-bench/src/jsonbench/data.rs +++ b/vortex-bench/src/jsonbench/data.rs @@ -63,7 +63,8 @@ const SHREDDING_MAX_DEPTH: usize = 6; /// URL of the `file_idx`-th (1-based) raw JSONBench file. pub fn raw_json_url(file_idx: usize) -> String { format!( - "https://clickhouse-public-datasets.s3.amazonaws.com/bluesky/file_{file_idx:04}.json.gz" + "https://clickhouse-public-datasets.s3.amazonaws.com/bluesky/{}", + raw_json_name(file_idx) ) } diff --git a/vortex-bench/src/jsonbench/mod.rs b/vortex-bench/src/jsonbench/mod.rs index eb432d1a8b4..d469b24f78c 100644 --- a/vortex-bench/src/jsonbench/mod.rs +++ b/vortex-bench/src/jsonbench/mod.rs @@ -7,7 +7,7 @@ //! Every event lives in a single `data` column of the `bluesky` table. The baseline Parquet format //! stores it as a JSON string; `parquet-variant` and the Vortex formats store it as a shredded //! Variant. The queries in `sql/jsonbench.sql` address JSON paths with placeholders that -//! [`JsonBenchBenchmark::query_for`] expands into each engine's idiom for each format. +//! [`expand_placeholders`] expands into each engine's idiom for each format. pub mod data; @@ -55,6 +55,7 @@ static PLACEHOLDER: LazyLock = LazyLock::new(|| { pub struct JsonBenchBenchmark { /// Number of one-million-row raw files the dataset is built from. n_files: usize, + data_dir: PathBuf, data_url: Url, /// Directory holding the downloaded raw files, shared by every scale. raw_dir: PathBuf, @@ -70,17 +71,12 @@ impl JsonBenchBenchmark { .map_err(|_| anyhow::anyhow!("invalid data directory {}", data_dir.display()))?; Ok(Self { n_files: scale_factor, + data_dir, data_url, raw_dir: root.join("json"), }) } - fn base_path(&self) -> anyhow::Result { - self.data_url - .to_file_path() - .map_err(|_| anyhow::anyhow!("jsonbench data URL must be a file URL")) - } - fn raw_path(&self, file_idx: usize) -> PathBuf { self.raw_dir.join(raw_json_name(file_idx)) } @@ -90,7 +86,7 @@ impl JsonBenchBenchmark { } fn shredding_schema(&self) -> anyhow::Result { - let path = self.base_path()?.join("shredding.json"); + let path = self.data_dir.join("shredding.json"); idempotent(&path, |tmp| { let schema = ShreddingSchema::infer(&self.raw_path(1))?; info!( @@ -146,7 +142,7 @@ impl JsonBenchBenchmark { fn prepare_variant_parquet(&self) -> anyhow::Result { let shredding = self.shredding_schema()?; - let dir = self.base_path()?.join(Format::ParquetVariant.name()); + let dir = self.data_dir.join(Format::ParquetVariant.name()); self.convert_files(&dir, "parquet", |raw, output| { write_variant_parquet(raw, output, &shredding) })?; @@ -158,8 +154,8 @@ impl JsonBenchBenchmark { Format::VortexCompact => CompactionStrategy::Compact, _ => CompactionStrategy::Default, }; - let variant_dir = self.prepare_variant_parquet()?; - let vortex_dir = self.base_path()?.join(format.name()); + let variant_dir = tokio::task::block_in_place(|| self.prepare_variant_parquet())?; + let vortex_dir = self.data_dir.join(format.name()); fs::create_dir_all(&vortex_dir)?; for file_idx in self.file_indices() { let stem = output_stem(file_idx); @@ -171,18 +167,16 @@ impl JsonBenchBenchmark { } Ok(()) } +} - /// Expand `{str:path}` and `{i64:path}` placeholders into `engine`'s idiom for reading that - /// JSON path from `format`. - pub fn expand_placeholders(engine: Engine, format: Format, query: &str) -> String { - PLACEHOLDER - .replace_all(query, |caps: &Captures| { - let ty = &caps[1]; - let path = &caps[2]; - expand_path(engine, format, ty, path) - }) - .into_owned() - } +/// Expand `{str:path}` and `{i64:path}` placeholders into `engine`'s idiom for reading that JSON +/// path from `format`. +pub fn expand_placeholders(engine: Engine, format: Format, query: &str) -> String { + PLACEHOLDER + .replace_all(query, |caps: &Captures| { + expand_path(engine, format, &caps[1], &caps[2]) + }) + .into_owned() } /// The SQL expression reading JSON `path` as `ty` (`str` or `i64`) on `engine` from `format`. @@ -240,7 +234,7 @@ impl Benchmark for JsonBenchBenchmark { } fn query_for(&self, engine: Engine, format: Format, query: &str) -> String { - Self::expand_placeholders(engine, format, query) + expand_placeholders(engine, format, query) } async fn generate_base_data(&self) -> anyhow::Result<()> { @@ -267,7 +261,7 @@ impl Benchmark for JsonBenchBenchmark { .await?; } - let parquet_dir = self.base_path()?.join(Format::Parquet.name()); + let parquet_dir = self.data_dir.join(Format::Parquet.name()); tokio::task::block_in_place(|| { self.shredding_schema()?; self.convert_files(&parquet_dir, "parquet", write_json_parquet) @@ -279,10 +273,7 @@ impl Benchmark for JsonBenchBenchmark { Format::ParquetVariant => { tokio::task::block_in_place(|| self.prepare_variant_parquet())?; } - Format::OnDiskVortex | Format::VortexCompact => { - tokio::task::block_in_place(|| self.prepare_variant_parquet())?; - self.prepare_vortex(format).await?; - } + Format::OnDiskVortex | Format::VortexCompact => self.prepare_vortex(format).await?, _ => {} } Ok(()) @@ -349,9 +340,6 @@ mod tests { #[case] format: Format, #[case] expected: &str, ) { - assert_eq!( - JsonBenchBenchmark::expand_placeholders(engine, format, QUERY), - expected - ); + assert_eq!(expand_placeholders(engine, format, QUERY), expected); } } From 1d6ff339a82b6caa916724fab753b9ad196a6df8 Mon Sep 17 00:00:00 2001 From: Joe Isaacs Date: Wed, 7 Oct 2026 12:40:48 +0000 Subject: [PATCH 3/3] Sort the workspace dependencies the way taplo expects Signed-off-by: Joe Isaacs Co-Authored-By: Claude Claude-Session: https://claude.ai/code/session_01P9G2RwAZnBNHUackdtcjeq --- Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Cargo.toml b/Cargo.toml index 509aeabac8d..2a4e0e515f1 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -162,8 +162,8 @@ env_logger = "0.11" fastlanes = { version = "0.7.0", features = ["runtime"] } fearless_simd = "1.0.0" fearless_simd_macros = "0.1.0" -flate2 = "1.1" flatbuffers = "25.2.10" +flate2 = "1.1" fsst-rs = "0.6.0" futures = { version = "0.3.31", default-features = false } fuzzy-matcher = "0.3"