From 182c8ac0deae182f6f31a297ff53f42d804b12df Mon Sep 17 00:00:00 2001 From: You-Cheng Lin Date: Sun, 27 Sep 2026 17:25:52 +0800 Subject: [PATCH 1/4] add q16 q19 to tpch test suite Signed-off-by: You-Cheng Lin --- .../nightly_tests/dataset/tpch/tpch_q16.py | 104 ++++++++++++++++++ .../nightly_tests/dataset/tpch/tpch_q19.py | 94 ++++++++++++++++ 2 files changed, 198 insertions(+) create mode 100644 release/nightly_tests/dataset/tpch/tpch_q16.py create mode 100644 release/nightly_tests/dataset/tpch/tpch_q19.py 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..42b63beffb9 --- /dev/null +++ b/release/nightly_tests/dataset/tpch/tpch_q16.py @@ -0,0 +1,104 @@ +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; + # + # Note: + # The NOT IN subquery is a left_anti join (as in Q22). COUNT(DISTINCT) + # is expressed as two groupbys: dedupe on (brand, type, size, suppkey), + # then count rows per (brand, type, 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"]) + + # NOT IN -> anti join. + ps_clean = partsupp.join( + complainers, + join_type="left_anti", + num_partitions=200, + on=("ps_suppkey",), + right_on=("s_suppkey",), + ) + + joined = ps_clean.join( + part_filtered, + join_type="inner", + num_partitions=200, + on=("ps_partkey",), + right_on=("p_partkey",), + ).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..36ba57b4ebe --- /dev/null +++ b/release/nightly_tests/dataset/tpch/tpch_q19.py @@ -0,0 +1,94 @@ +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 ...); + # + # Note: + # The shipmode/shipinstruct predicates are common to all three clauses + # and filter lineitem before the join; the disjunction of the remaining + # brand/container/quantity/size conjunctions runs after the join. + + 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), + ] + + lineitem_filtered = lineitem.filter( + expr=col("l_shipmode").is_in(["AIR", "AIR REG"]) + & (col("l_shipinstruct") == "DELIVER IN PERSON") + ) + + joined = lineitem_filtered.join( + part, + join_type="inner", + num_partitions=200, + on=("l_partkey",), + right_on=("p_partkey",), + ) + + disjunction = None + for brand, containers, qty_lo, qty_hi, size_hi in clauses: + clause = ( + (col("p_brand") == brand) + & col("p_container").is_in(containers) + & (col("l_quantity") >= qty_lo) + & (col("l_quantity") <= qty_hi) + & (col("p_size") >= 1) + & (col("p_size") <= size_hi) + ) + disjunction = clause if disjunction is None else (disjunction | clause) + + 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) From b0e3f141d3811ce39f5da0f55d4b3f07d4034bbb Mon Sep 17 00:00:00 2001 From: You-Cheng Lin Date: Sun, 27 Sep 2026 17:28:04 +0800 Subject: [PATCH 2/4] add yaml changes Signed-off-by: You-Cheng Lin --- release/release_data_tests.yaml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 From 0544dc7566646ab0e226d37ff4ffdfde38839227 Mon Sep 17 00:00:00 2001 From: You-Cheng Lin Date: Sun, 27 Sep 2026 17:54:38 +0800 Subject: [PATCH 3/4] update benchmark query Signed-off-by: You-Cheng Lin --- .../nightly_tests/dataset/tpch/tpch_q16.py | 20 ++++---- .../nightly_tests/dataset/tpch/tpch_q19.py | 49 +++++++++++++------ 2 files changed, 44 insertions(+), 25 deletions(-) diff --git a/release/nightly_tests/dataset/tpch/tpch_q16.py b/release/nightly_tests/dataset/tpch/tpch_q16.py index 42b63beffb9..95383bfb521 100644 --- a/release/nightly_tests/dataset/tpch/tpch_q16.py +++ b/release/nightly_tests/dataset/tpch/tpch_q16.py @@ -58,21 +58,23 @@ def benchmark_fn(): 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. - ps_clean = partsupp.join( + joined = ps_filtered.join( complainers, join_type="left_anti", num_partitions=200, on=("ps_suppkey",), right_on=("s_suppkey",), - ) - - joined = ps_clean.join( - part_filtered, - join_type="inner", - num_partitions=200, - on=("ps_partkey",), - right_on=("p_partkey",), ).select_columns(["p_brand", "p_type", "p_size", "ps_suppkey"]) # COUNT(DISTINCT ps_suppkey): dedupe first, then count. diff --git a/release/nightly_tests/dataset/tpch/tpch_q19.py b/release/nightly_tests/dataset/tpch/tpch_q19.py index 36ba57b4ebe..2b022e99b81 100644 --- a/release/nightly_tests/dataset/tpch/tpch_q19.py +++ b/release/nightly_tests/dataset/tpch/tpch_q19.py @@ -23,9 +23,11 @@ def benchmark_fn(): # OR (... 'Brand#34', LG containers, quantity 20..30, size 1..15 ...); # # Note: - # The shipmode/shipinstruct predicates are common to all three clauses - # and filter lineitem before the join; the disjunction of the remaining - # brand/container/quantity/size conjunctions runs after the join. + # Weaker implied predicates are pushed below the join: the disjunction + # of the part-only (brand/container/size) conjunctions onto part, and + # the shared shipmode/shipinstruct predicates plus the global quantity + # bounds onto lineitem. The exact per-clause disjunction, which ties + # each brand to its quantity range, still runs after the join. part = load_table("part", args.sf).select_columns( ["p_partkey", "p_brand", "p_size", "p_container"] @@ -48,31 +50,46 @@ def benchmark_fn(): ("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, + part_filtered, join_type="inner", num_partitions=200, on=("l_partkey",), right_on=("p_partkey",), ) - disjunction = None - for brand, containers, qty_lo, qty_hi, size_hi in clauses: - clause = ( - (col("p_brand") == brand) - & col("p_container").is_in(containers) - & (col("l_quantity") >= qty_lo) - & (col("l_quantity") <= qty_hi) - & (col("p_size") >= 1) - & (col("p_size") <= size_hi) - ) - disjunction = clause if disjunction is None else (disjunction | clause) - ds = joined.filter(expr=disjunction) ds = ds.with_column( From d8dc50b1ca59f05fccfbd2dcb2196f104a8f6816 Mon Sep 17 00:00:00 2001 From: You-Cheng Lin Date: Sun, 27 Sep 2026 18:06:43 +0800 Subject: [PATCH 4/4] remove note comments Signed-off-by: You-Cheng Lin --- release/nightly_tests/dataset/tpch/tpch_q16.py | 5 ----- release/nightly_tests/dataset/tpch/tpch_q19.py | 7 ------- 2 files changed, 12 deletions(-) diff --git a/release/nightly_tests/dataset/tpch/tpch_q16.py b/release/nightly_tests/dataset/tpch/tpch_q16.py index 95383bfb521..85e652f8802 100644 --- a/release/nightly_tests/dataset/tpch/tpch_q16.py +++ b/release/nightly_tests/dataset/tpch/tpch_q16.py @@ -23,11 +23,6 @@ def benchmark_fn(): # 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; - # - # Note: - # The NOT IN subquery is a left_anti join (as in Q22). COUNT(DISTINCT) - # is expressed as two groupbys: dedupe on (brand, type, size, suppkey), - # then count rows per (brand, type, size). part = load_table("part", args.sf).select_columns( ["p_partkey", "p_brand", "p_type", "p_size"] diff --git a/release/nightly_tests/dataset/tpch/tpch_q19.py b/release/nightly_tests/dataset/tpch/tpch_q19.py index 2b022e99b81..d832d7825c8 100644 --- a/release/nightly_tests/dataset/tpch/tpch_q19.py +++ b/release/nightly_tests/dataset/tpch/tpch_q19.py @@ -21,13 +21,6 @@ def benchmark_fn(): # 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 ...); - # - # Note: - # Weaker implied predicates are pushed below the join: the disjunction - # of the part-only (brand/container/size) conjunctions onto part, and - # the shared shipmode/shipinstruct predicates plus the global quantity - # bounds onto lineitem. The exact per-clause disjunction, which ties - # each brand to its quantity range, still runs after the join. part = load_table("part", args.sf).select_columns( ["p_partkey", "p_brand", "p_size", "p_container"]