Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 37 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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" }
Expand All @@ -162,6 +163,7 @@ fastlanes = { version = "0.7.0", features = ["runtime"] }
fearless_simd = "1.0.0"
fearless_simd_macros = "0.1.0"
flatbuffers = "25.2.10"
flate2 = "1.1"
fsst-rs = "0.6.0"
futures = { version = "0.3.31", default-features = false }
fuzzy-matcher = "0.3"
Expand Down
2 changes: 1 addition & 1 deletion bench-orchestrator/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ vx-bench run <benchmark> [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:**

Expand Down
1 change: 1 addition & 0 deletions bench-orchestrator/bench_orchestrator/ci_matrix/render.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
4 changes: 4 additions & 0 deletions bench-orchestrator/bench_orchestrator/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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"
Expand All @@ -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,
Expand Down
1 change: 1 addition & 0 deletions benchmarks/datafusion-bench/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
16 changes: 15 additions & 1 deletion benchmarks/datafusion-bench/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -101,7 +115,7 @@ pub fn make_object_store(
pub fn format_to_df_format(format: Format) -> anyhow::Result<Arc<dyn FileFormat>> {
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()),
),
Expand Down
6 changes: 4 additions & 2 deletions benchmarks/datafusion-bench/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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?;
}
Expand All @@ -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();
Expand Down
1 change: 1 addition & 0 deletions benchmarks/duckdb-bench/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -208,6 +208,7 @@ impl DuckClient {

let object_type = match file_format {
Format::Parquet
| Format::ParquetVariant
| Format::OnDiskVortex
| Format::VortexCompact
| Format::VortexSpatialNative => "VIEW",
Expand Down
6 changes: 4 additions & 2 deletions benchmarks/duckdb-bench/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -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(())
Expand Down Expand Up @@ -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))
},
)?;

Expand Down
9 changes: 8 additions & 1 deletion vortex-bench/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand All @@ -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"] }
Expand Down
66 changes: 66 additions & 0 deletions vortex-bench/sql/jsonbench.md
Original file line number Diff line number Diff line change
@@ -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.
49 changes: 49 additions & 0 deletions vortex-bench/sql/jsonbench.sql
Original file line number Diff line number Diff line change
@@ -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;
Loading
Loading