Skip to content
Merged
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
101 changes: 101 additions & 0 deletions release/nightly_tests/dataset/tpch/tpch_q16.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
import ray
from ray.data.aggregate import Count
from ray.data.expressions import col
from common import load_table, parse_tpch_args, run_tpch_benchmark


def main(args):
def benchmark_fn():
# Q16: Parts/Supplier Relationship Query
# Count distinct suppliers able to supply parts of given sizes,
# excluding a brand, a type prefix, and suppliers with complaints.
#
# Equivalent SQL:
# SELECT p_brand, p_type, p_size,
# COUNT(DISTINCT ps_suppkey) AS supplier_cnt
# FROM partsupp, part
# WHERE p_partkey = ps_partkey
# AND p_brand <> 'Brand#45'
# AND p_type NOT LIKE 'MEDIUM POLISHED%'
# AND p_size IN (49, 14, 23, 45, 19, 3, 36, 9)
# AND ps_suppkey NOT IN (
# SELECT s_suppkey FROM supplier
# WHERE s_comment LIKE '%Customer%Complaints%')
# GROUP BY p_brand, p_type, p_size
# ORDER BY supplier_cnt DESC, p_brand, p_type, p_size;

part = load_table("part", args.sf).select_columns(
["p_partkey", "p_brand", "p_type", "p_size"]
)
partsupp = load_table("partsupp", args.sf).select_columns(
["ps_partkey", "ps_suppkey"]
)
supplier = load_table("supplier", args.sf).select_columns(
["s_suppkey", "s_comment"]
)

# Q16 parameters
excluded_brand = "Brand#45"
excluded_type_prefix = "MEDIUM POLISHED"
sizes = [49, 14, 23, 45, 19, 3, 36, 9]

part_filtered = part.filter(
expr=(col("p_brand") != excluded_brand) & col("p_size").is_in(sizes)
)
# NOT LIKE 'MEDIUM POLISHED%'. Kept after load_table to avoid pushing
# a UDF expression into parquet predicate conversion (see Q2/Q9).
part_filtered = part_filtered.filter(
expr=~col("p_type").str.starts_with(excluded_type_prefix)
)

# Suppliers with complaints: LIKE '%Customer%Complaints%'.
complainers = supplier.filter(
expr=col("s_comment").str.match_regex("Customer.*Complaints")
).select_columns(["s_suppkey"])

# Inner join with the selective part filter first so the anti join
# below shuffles the reduced dataset instead of all of partsupp.
ps_filtered = partsupp.join(
part_filtered,
join_type="inner",
num_partitions=200,
on=("ps_partkey",),
right_on=("p_partkey",),
)

# NOT IN -> anti join.
joined = ps_filtered.join(
complainers,
join_type="left_anti",
num_partitions=200,
on=("ps_suppkey",),
right_on=("s_suppkey",),
).select_columns(["p_brand", "p_type", "p_size", "ps_suppkey"])
Comment thread
owenowenisme marked this conversation as resolved.

# COUNT(DISTINCT ps_suppkey): dedupe first, then count.
distinct_suppliers = (
joined.groupby(["p_brand", "p_type", "p_size", "ps_suppkey"])
.aggregate(Count(alias_name="_dedupe"))
.select_columns(["p_brand", "p_type", "p_size", "ps_suppkey"])
)

_ = (
distinct_suppliers.groupby(["p_brand", "p_type", "p_size"])
.aggregate(Count(alias_name="supplier_cnt"))
.sort(
key=["supplier_cnt", "p_brand", "p_type", "p_size"],
descending=[True, False, False, False],
)
.materialize()
)

# Report arguments for the benchmark.
return vars(args)

run_tpch_benchmark("tpch_q16", benchmark_fn)


if __name__ == "__main__":
ray.init()
args = parse_tpch_args()
main(args)
104 changes: 104 additions & 0 deletions release/nightly_tests/dataset/tpch/tpch_q19.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
import ray
from ray.data.aggregate import Sum
from ray.data.expressions import col
from common import load_table, parse_tpch_args, run_tpch_benchmark, to_f64


def main(args):
def benchmark_fn():
# Q19: Discounted Revenue Query
# Revenue from qualifying part/lineitem combinations across three
# brand/container/quantity/size clauses.
#
# Equivalent SQL:
# SELECT SUM(l_extendedprice * (1 - l_discount)) AS revenue
# FROM lineitem, part
# WHERE (p_partkey = l_partkey AND p_brand = 'Brand#12'
# AND p_container IN ('SM CASE','SM BOX','SM PACK','SM PKG')
# AND l_quantity >= 1 AND l_quantity <= 11
# AND p_size BETWEEN 1 AND 5
# AND l_shipmode IN ('AIR','AIR REG')
# AND l_shipinstruct = 'DELIVER IN PERSON')
# OR (... 'Brand#23', MED containers, quantity 10..20, size 1..10 ...)
# OR (... 'Brand#34', LG containers, quantity 20..30, size 1..15 ...);

part = load_table("part", args.sf).select_columns(
["p_partkey", "p_brand", "p_size", "p_container"]
)
lineitem = load_table("lineitem", args.sf).select_columns(
[
"l_partkey",
"l_quantity",
"l_extendedprice",
"l_discount",
"l_shipinstruct",
"l_shipmode",
]
)

# Q19 parameters
clauses = [
("Brand#12", ["SM CASE", "SM BOX", "SM PACK", "SM PKG"], 1, 11, 5),
("Brand#23", ["MED BAG", "MED BOX", "MED PKG", "MED PACK"], 10, 20, 10),
("Brand#34", ["LG CASE", "LG BOX", "LG PACK", "LG PKG"], 20, 30, 15),
]

part_disjunction = None
disjunction = None
for brand, containers, qty_lo, qty_hi, size_hi in clauses:
part_clause = (
(col("p_brand") == brand)
& col("p_container").is_in(containers)
& (col("p_size") >= 1)
& (col("p_size") <= size_hi)
)
clause = (
part_clause
& (col("l_quantity") >= qty_lo)
& (col("l_quantity") <= qty_hi)
)
part_disjunction = (
part_clause
if part_disjunction is None
else (part_disjunction | part_clause)
)
disjunction = clause if disjunction is None else (disjunction | clause)

part_filtered = part.filter(expr=part_disjunction)

min_qty = min(qty_lo for _, _, qty_lo, _, _ in clauses)
max_qty = max(qty_hi for _, _, _, qty_hi, _ in clauses)
lineitem_filtered = lineitem.filter(
expr=col("l_shipmode").is_in(["AIR", "AIR REG"])
& (col("l_shipinstruct") == "DELIVER IN PERSON")
& (col("l_quantity") >= min_qty)
& (col("l_quantity") <= max_qty)
)

joined = lineitem_filtered.join(
part_filtered,
join_type="inner",
num_partitions=200,
on=("l_partkey",),
right_on=("p_partkey",),
)
Comment thread
owenowenisme marked this conversation as resolved.

ds = joined.filter(expr=disjunction)

ds = ds.with_column(
"revenue",
to_f64(col("l_extendedprice")) * (1 - to_f64(col("l_discount"))),
)

_ = ds.aggregate(Sum(on="revenue", alias_name="revenue"))

# Report arguments for the benchmark.
return vars(args)

run_tpch_benchmark("tpch_q19", benchmark_fn)


if __name__ == "__main__":
ray.init()
args = parse_tpch_args()
main(args)
2 changes: 1 addition & 1 deletion release/release_data_tests.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -1255,7 +1255,7 @@
frequency: manual
matrix:
setup:
query: [q2, q3, q4, q5, q6, q7, q9, q10, q11, q12, q14, q17, q18, q20]
query: [q2, q3, q4, q5, q6, q7, q9, q10, q11, q12, q14, q16, q17, q18, q19, q20]

cluster:
anyscale_sdk_2026: true
Expand Down
Loading