diff --git a/release/nightly_tests/dataset/tpch/tpch_q16.py b/release/nightly_tests/dataset/tpch/tpch_q16.py new file mode 100644 index 00000000000..85e652f8802 --- /dev/null +++ b/release/nightly_tests/dataset/tpch/tpch_q16.py @@ -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"]) + + # 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) diff --git a/release/nightly_tests/dataset/tpch/tpch_q19.py b/release/nightly_tests/dataset/tpch/tpch_q19.py new file mode 100644 index 00000000000..d832d7825c8 --- /dev/null +++ b/release/nightly_tests/dataset/tpch/tpch_q19.py @@ -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",), + ) + + 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) diff --git a/release/release_data_tests.yaml b/release/release_data_tests.yaml index fff77267dc7..8bbc0520fa9 100644 --- a/release/release_data_tests.yaml +++ b/release/release_data_tests.yaml @@ -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