Skip to content

Commit 3fd6879

Browse files
singaraionaclaude
andcommitted
Add sort and window join benchmarks
- Add sort_single and sort_multi tasks to all adapters - Add window_join task to all adapters (wj1 equivalent) - Create sort.yaml benchmark suite (6 queries) - Create window_join.yaml benchmark suite - Add h2oai_join_1e7 and window_join_10m dataset manifests - Add generate_window_join.py script for data generation Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
1 parent af8f783 commit 3fd6879

9 files changed

Lines changed: 501 additions & 10 deletions

File tree

adapters/duckdb_adapter.py

Lines changed: 60 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,11 @@ def _build_task_registry(self) -> dict[str, callable]:
7777
# Join queries
7878
"inner_join": self._task_inner_join,
7979
"left_join": self._task_left_join,
80+
# Sort queries
81+
"sort_single": self._task_sort_single,
82+
"sort_multi": self._task_sort_multi,
83+
# Window join queries
84+
"window_join": self._task_window_join,
8085
# Generic SQL execution
8186
"sql": self._task_sql,
8287
}
@@ -307,13 +312,64 @@ def _task_left_join(self, params: dict[str, Any]) -> AdapterResult:
307312
left_table = params.get("left_table", "x")
308313
right_table = params.get("right_table", "y")
309314
query = f"""
310-
SELECT * FROM {left_table}
311-
LEFT JOIN {right_table}
312-
ON {left_table}.id1 = {right_table}.id1
315+
SELECT * FROM {left_table}
316+
LEFT JOIN {right_table}
317+
ON {left_table}.id1 = {right_table}.id1
313318
AND {left_table}.id2 = {right_table}.id2
314319
"""
315320
return self._execute_query(query)
316-
321+
322+
# =========================================================================
323+
# Sort Queries (SQL syntax)
324+
# =========================================================================
325+
326+
def _task_sort_single(self, params: dict[str, Any]) -> AdapterResult:
327+
"""Sort by single column"""
328+
table = params.get("table", self._table_name)
329+
column = params.get("column", "id1")
330+
descending = params.get("descending", False)
331+
order = "DESC" if descending else "ASC"
332+
query = f"SELECT * FROM {table} ORDER BY {column} {order}"
333+
return self._execute_query(query)
334+
335+
def _task_sort_multi(self, params: dict[str, Any]) -> AdapterResult:
336+
"""Sort by multiple columns"""
337+
table = params.get("table", self._table_name)
338+
columns = params.get("columns", ["id1", "id2"])
339+
order_by = ", ".join(columns)
340+
query = f"SELECT * FROM {table} ORDER BY {order_by}"
341+
return self._execute_query(query)
342+
343+
# =========================================================================
344+
# Window Join Queries (SQL ASOF JOIN)
345+
# DuckDB supports ASOF JOIN for time-series window joins
346+
# =========================================================================
347+
348+
def _task_window_join(self, params: dict[str, Any]) -> AdapterResult:
349+
"""Window join using ASOF JOIN - join within time window with aggregations
350+
351+
Note: DuckDB's ASOF JOIN is point-in-time, not a true window aggregation.
352+
For a fair comparison, we use a range join with aggregation.
353+
"""
354+
trades_table = params.get("trades_table", "trades")
355+
quotes_table = params.get("quotes_table", "quotes")
356+
window_ms = params.get("window_ms", 10000) # +/- 10 seconds default
357+
358+
# Use range join with aggregation for true window join semantics
359+
query = f"""
360+
SELECT
361+
t.*,
362+
MIN(q.Bid) as Bid,
363+
MAX(q.Ask) as Ask
364+
FROM {trades_table} t
365+
LEFT JOIN {quotes_table} q
366+
ON t.Sym = q.Sym
367+
AND q.Ts BETWEEN t.Ts - INTERVAL '{window_ms} milliseconds'
368+
AND t.Ts + INTERVAL '{window_ms} milliseconds'
369+
GROUP BY ALL
370+
"""
371+
return self._execute_query(query)
372+
317373
def _task_sql(self, params: dict[str, Any]) -> AdapterResult:
318374
"""Execute arbitrary SQL query."""
319375
query = params.get("query")

adapters/kdb_adapter.py

Lines changed: 47 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,11 @@ def _build_task_registry(self) -> dict[str, callable]:
7070
# Join queries
7171
"inner_join": self._task_inner_join,
7272
"left_join": self._task_left_join,
73+
# Sort queries
74+
"sort_single": self._task_sort_single,
75+
"sort_multi": self._task_sort_multi,
76+
# Window join queries
77+
"window_join": self._task_window_join,
7378
# Generic q execution
7479
"eval": self._task_eval,
7580
}
@@ -343,7 +348,48 @@ def _task_left_join(self, params: dict[str, Any]) -> AdapterResult:
343348
# lj = left join, need to key the right table first
344349
expr = f"{left_table} lj `id1`id2 xkey {right_table}"
345350
return self._execute_q(expr)
346-
351+
352+
# =========================================================================
353+
# Sort Queries (q syntax)
354+
# xasc = ascending sort, xdesc = descending sort
355+
# =========================================================================
356+
357+
def _task_sort_single(self, params: dict[str, Any]) -> AdapterResult:
358+
"""Sort by single column"""
359+
table = params.get("table", self._table_name)
360+
column = params.get("column", "id1")
361+
descending = params.get("descending", False)
362+
if descending:
363+
expr = f"`{column} xdesc {table}"
364+
else:
365+
expr = f"`{column} xasc {table}"
366+
return self._execute_q(expr)
367+
368+
def _task_sort_multi(self, params: dict[str, Any]) -> AdapterResult:
369+
"""Sort by multiple columns"""
370+
table = params.get("table", self._table_name)
371+
columns = params.get("columns", ["id1", "id2"])
372+
# q multi-column sort: `col1`col2 xasc t
373+
cols_str = "`" + "`".join(columns)
374+
expr = f"{cols_str} xasc {table}"
375+
return self._execute_q(expr)
376+
377+
# =========================================================================
378+
# Window Join Queries (q syntax)
379+
# wj1 = window join with prevailing values
380+
# =========================================================================
381+
382+
def _task_window_join(self, params: dict[str, Any]) -> AdapterResult:
383+
"""Window join (wj1) - join within time window with aggregations"""
384+
trades_table = params.get("trades_table", "trades")
385+
quotes_table = params.get("quotes_table", "quotes")
386+
window_ms = params.get("window_ms", 10000) # +/- 10 seconds default
387+
388+
# q window join: wj1[w;`Sym`Ts;trades;(quotes;(min;`Bid);(max;`Ask))]
389+
# w is a 2-row matrix of [start_times; end_times]
390+
expr = f"wj1[(-{window_ms};{window_ms})+\\:{trades_table}.Ts;`Sym`Ts;{trades_table};(`Sym`Ts xasc {quotes_table};(min;`Bid);(max;`Ask))]"
391+
return self._execute_q(expr)
392+
347393
def _task_eval(self, params: dict[str, Any]) -> AdapterResult:
348394
"""Execute arbitrary q expression."""
349395
expr = params.get("expr")

adapters/polars_adapter.py

Lines changed: 81 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -98,6 +98,11 @@ def _build_task_registry(self) -> dict[str, callable]:
9898
# Join queries
9999
"inner_join": self._task_inner_join,
100100
"left_join": self._task_left_join,
101+
# Sort queries
102+
"sort_single": self._task_sort_single,
103+
"sort_multi": self._task_sort_multi,
104+
# Window join queries
105+
"window_join": self._task_window_join,
101106
# Generic execution
102107
"sql": self._task_sql,
103108
}
@@ -468,15 +473,87 @@ def _task_left_join(self, params: dict[str, Any]) -> AdapterResult:
468473
"""Left join on id1, id2"""
469474
left_table = params.get("left_table", "x")
470475
right_table = params.get("right_table", "y")
471-
476+
472477
left_lf = self._get_lazy_table(left_table)
473478
right_lf = self._get_lazy_table(right_table)
474-
479+
475480
def query():
476481
return left_lf.join(right_lf, on=["id1", "id2"], how="left")
477-
482+
478483
return self._execute_lazy(query, f'{left_table}.join({right_table}, on=["id1","id2"], how="left")')
479-
484+
485+
# =========================================================================
486+
# Sort Queries (using lazy evaluation for optimal parallelism)
487+
# =========================================================================
488+
489+
def _task_sort_single(self, params: dict[str, Any]) -> AdapterResult:
490+
"""Sort by single column"""
491+
table_name = params.get("table", self._table_name)
492+
column = params.get("column", "id1")
493+
descending = params.get("descending", False)
494+
lf = self._get_lazy_table(table_name)
495+
496+
def query():
497+
return lf.sort(column, descending=descending)
498+
499+
return self._execute_lazy(query, f'lf.sort("{column}", descending={descending})')
500+
501+
def _task_sort_multi(self, params: dict[str, Any]) -> AdapterResult:
502+
"""Sort by multiple columns"""
503+
table_name = params.get("table", self._table_name)
504+
columns = params.get("columns", ["id1", "id2"])
505+
lf = self._get_lazy_table(table_name)
506+
507+
def query():
508+
return lf.sort(columns)
509+
510+
return self._execute_lazy(query, f'lf.sort({columns})')
511+
512+
# =========================================================================
513+
# Window Join Queries (using join_asof for time-series joins)
514+
# =========================================================================
515+
516+
def _task_window_join(self, params: dict[str, Any]) -> AdapterResult:
517+
"""Window join - join within time window with aggregations
518+
519+
Polars doesn't have a direct wj1 equivalent, so we use a range join
520+
approach with group_by for aggregation.
521+
"""
522+
import datetime
523+
524+
trades_table = params.get("trades_table", "trades")
525+
quotes_table = params.get("quotes_table", "quotes")
526+
window_ms = params.get("window_ms", 10000) # +/- 10 seconds default
527+
528+
trades_lf = self._get_lazy_table(trades_table)
529+
quotes_lf = self._get_lazy_table(quotes_table)
530+
531+
def query():
532+
# Add window boundaries to trades
533+
trades_with_window = trades_lf.with_columns([
534+
(pl.col("Ts") - datetime.timedelta(milliseconds=window_ms)).alias("_window_start"),
535+
(pl.col("Ts") + datetime.timedelta(milliseconds=window_ms)).alias("_window_end"),
536+
])
537+
538+
# Cross join on Sym, filter by time window, then aggregate
539+
# This is expensive but semantically correct
540+
result = (
541+
trades_with_window
542+
.join(quotes_lf, on="Sym", how="left")
543+
.filter(
544+
(pl.col("Ts_right") >= pl.col("_window_start")) &
545+
(pl.col("Ts_right") <= pl.col("_window_end"))
546+
)
547+
.group_by(["Sym", "Ts", "Price"])
548+
.agg([
549+
pl.min("Bid").alias("Bid"),
550+
pl.max("Ask").alias("Ask"),
551+
])
552+
)
553+
return result
554+
555+
return self._execute_lazy(query, 'window_join(trades, quotes)')
556+
480557
def _task_sql(self, params: dict[str, Any]) -> AdapterResult:
481558
"""Execute arbitrary SQL query."""
482559
query = params.get("query")

adapters/rayforce_adapter.py

Lines changed: 50 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -195,6 +195,11 @@ def _build_task_registry(self) -> dict[str, callable]:
195195
# Join queries
196196
"inner_join": self._task_inner_join,
197197
"left_join": self._task_left_join,
198+
# Sort queries
199+
"sort_single": self._task_sort_single,
200+
"sort_multi": self._task_sort_multi,
201+
# Window join queries
202+
"window_join": self._task_window_join,
198203
# Generic expression execution
199204
"eval": self._task_eval,
200205
}
@@ -612,7 +617,51 @@ def _task_left_join(self, params: dict[str, Any]) -> AdapterResult:
612617
right_table = params.get("right_table", "y")
613618
expr = f"(left-join [id1 id2] {left_table} {right_table})"
614619
return self._execute_expr(expr)
615-
620+
621+
# =========================================================================
622+
# Sort Queries (Rayforce syntax)
623+
# From Rayforce docs: (sort col t) or (asc col t)
624+
# =========================================================================
625+
626+
def _task_sort_single(self, params: dict[str, Any]) -> AdapterResult:
627+
"""Sort by single column"""
628+
table = params.get("table", self._table_name)
629+
column = params.get("column", "id1")
630+
descending = params.get("descending", False)
631+
if descending:
632+
expr = f"(desc {column} {table})"
633+
else:
634+
expr = f"(asc {column} {table})"
635+
return self._execute_expr(expr)
636+
637+
def _task_sort_multi(self, params: dict[str, Any]) -> AdapterResult:
638+
"""Sort by multiple columns"""
639+
table = params.get("table", self._table_name)
640+
columns = params.get("columns", ["id1", "id2"])
641+
# Rayforce multi-column sort: (asc [col1 col2] t)
642+
cols_str = " ".join(columns)
643+
expr = f"(asc [{cols_str}] {table})"
644+
return self._execute_expr(expr)
645+
646+
# =========================================================================
647+
# Window Join Queries (Rayforce syntax)
648+
# From Rayforce docs: (window-join1 [Sym Ts] intervals trades quotes {aggs})
649+
# =========================================================================
650+
651+
def _task_window_join(self, params: dict[str, Any]) -> AdapterResult:
652+
"""Window join (wj1) - join within time window with aggregations"""
653+
trades_table = params.get("trades_table", "trades")
654+
quotes_table = params.get("quotes_table", "quotes")
655+
keys = params.get("keys", ["Sym", "Ts"])
656+
window_ms = params.get("window_ms", 10000) # +/- 10 seconds default
657+
658+
keys_str = " ".join(keys)
659+
# Build intervals from trades timestamp
660+
expr = f"""(do
661+
(set _intervals (map-left + [-{window_ms} {window_ms}] (at {trades_table} 'Ts)))
662+
(window-join1 [{keys_str}] _intervals {trades_table} {quotes_table} {{Bid: (min Bid) Ask: (max Ask)}}))"""
663+
return self._execute_expr(expr)
664+
616665
def _task_eval(self, params: dict[str, Any]) -> AdapterResult:
617666
"""Execute arbitrary Rayforce expression."""
618667
expr = params.get("expr")
Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,43 @@
1+
{
2+
"name": "h2oai_join_1e7",
3+
"description": "H2OAI Join benchmark: 10M rows x table, join on id1, id2",
4+
"format": "h2oai_join",
5+
"version": "1.0.0",
6+
7+
"tables": {
8+
"x": {
9+
"table_name": "x",
10+
"row_count": 10000000,
11+
"columns": [
12+
{"name": "id1", "type": "I64"},
13+
{"name": "id2", "type": "I64"},
14+
{"name": "id3", "type": "I64"},
15+
{"name": "id4", "type": "SYMBOL"},
16+
{"name": "id5", "type": "SYMBOL"},
17+
{"name": "id6", "type": "SYMBOL"},
18+
{"name": "v1", "type": "F64"}
19+
],
20+
"files": ["J1_1e7_NA_0_0.csv"]
21+
},
22+
"y": {
23+
"table_name": "y",
24+
"row_count": 10000000,
25+
"columns": [
26+
{"name": "id1", "type": "I64"},
27+
{"name": "id2", "type": "I64"},
28+
{"name": "id3", "type": "I64"},
29+
{"name": "id4", "type": "SYMBOL"},
30+
{"name": "id5", "type": "SYMBOL"},
31+
{"name": "id6", "type": "SYMBOL"},
32+
{"name": "v2", "type": "F64"}
33+
],
34+
"files": ["J1_1e7_1e7_0_0.csv"]
35+
}
36+
},
37+
38+
"generation": {
39+
"tool": "Rscript _data/join-datagen.R 1e7 1e2 0 0",
40+
"source": "https://h2oai.github.io/db-benchmark",
41+
"description": "Standard H2OAI join benchmark with two 10M row tables"
42+
}
43+
}
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
{
2+
"name": "window_join_10m",
3+
"description": "Window join benchmark: 10M trades, 20M quotes",
4+
"format": "window_join",
5+
"version": "1.0.0",
6+
7+
"tables": {
8+
"trades": {
9+
"table_name": "trades",
10+
"row_count": 10000000,
11+
"columns": [
12+
{"name": "Sym", "type": "SYMBOL"},
13+
{"name": "Ts", "type": "TIME"},
14+
{"name": "Price", "type": "I64"}
15+
],
16+
"files": ["trades.csv"]
17+
},
18+
"quotes": {
19+
"table_name": "quotes",
20+
"row_count": 20000000,
21+
"columns": [
22+
{"name": "Sym", "type": "SYMBOL"},
23+
{"name": "Ts", "type": "TIME"},
24+
{"name": "Bid", "type": "I64"},
25+
{"name": "Ask", "type": "I64"}
26+
],
27+
"files": ["quotes.csv"]
28+
}
29+
},
30+
31+
"generation": {
32+
"tool": "scripts/generate_window_join.py",
33+
"description": "Trades: 99% AAPL, 1% MSFT. Quotes: AAPL/MSFT/GOOG distribution"
34+
}
35+
}

0 commit comments

Comments
 (0)