Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
26 commits
Select commit Hold shift + click to select a range
3e61e6b
Capture query plan in pds-h output
TomAugspurger Nov 7, 2025
aed6708
use string IDs
TomAugspurger Feb 3, 2026
1234b64
Merge remote-tracking branch 'upstream/main' into tom/structured-plans
TomAugspurger Feb 3, 2026
66e3356
fix
TomAugspurger Feb 3, 2026
01a42ce
fixes
TomAugspurger Feb 3, 2026
fd6a52c
test debugging
TomAugspurger Feb 4, 2026
748731f
polars compat
TomAugspurger Feb 4, 2026
0ebdbcc
Merge remote-tracking branch 'upstream/main' into tom/structured-plans
TomAugspurger Feb 4, 2026
53f30b1
cache stable id
TomAugspurger Feb 4, 2026
5924b82
test coverage
TomAugspurger Feb 4, 2026
1813d88
Merge remote-tracking branch 'upstream/main' into tom/structured-plans
TomAugspurger Feb 4, 2026
89f34dd
testing
TomAugspurger Feb 4, 2026
0429ef1
coverage
TomAugspurger Feb 4, 2026
fe6e794
doc strings
TomAugspurger Feb 5, 2026
1333359
Merge remote-tracking branch 'upstream/main' into tom/structured-plans
TomAugspurger Feb 5, 2026
6295afe
more docs
TomAugspurger Feb 5, 2026
bb57d6f
polars 1.30 compat
TomAugspurger Feb 5, 2026
71e7312
Merge remote-tracking branch 'upstream/main' into tom/structured-plans
TomAugspurger Feb 5, 2026
f3f908a
Move to get_stable_id
TomAugspurger Feb 5, 2026
a1ea778
Updates for logging changes
TomAugspurger Feb 5, 2026
e60fc0a
Merge branch 'main' into tom/structured-plans
TomAugspurger Feb 6, 2026
a8868f8
Merge remote-tracking branch 'upstream/main' into tom/structured-plans
TomAugspurger Feb 6, 2026
75027f5
dag -> plan
TomAugspurger Feb 6, 2026
a01d803
Merge remote-tracking branch 'upstream/main' into tom/structured-plans
TomAugspurger Feb 9, 2026
5728aa9
example
TomAugspurger Feb 9, 2026
a85e322
Merge branch 'main' into tom/structured-plans
TomAugspurger Feb 10, 2026
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
28 changes: 27 additions & 1 deletion python/cudf_polars/cudf_polars/dsl/nodebase.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@

from __future__ import annotations

import hashlib
from typing import TYPE_CHECKING, Any, ClassVar, Generic, TypeVar

if TYPE_CHECKING:
Expand Down Expand Up @@ -33,8 +34,9 @@ class Node(Generic[T]):
*children).``
"""

__slots__ = ("_hash_value", "_repr_value", "children")
__slots__ = ("_hash_value", "_repr_value", "_stable_hash_value", "children")
_hash_value: int
_stable_hash_value: int
_repr_value: str
children: tuple[T, ...]
_non_child: ClassVar[tuple[str, ...]] = ()
Expand Down Expand Up @@ -81,6 +83,30 @@ def get_hashable(self) -> Hashable:
"""
return (type(self), self._ctor_arguments(self.children))

def get_stable_id(self) -> int:
"""
Compute a stable identifier for Node.

Uses MD5 hash of the node's hashable representation for determinism
across process boundaries (Python's hash() uses PYTHONHASHSEED).

Parameters
----------
ir_node
The IR node.

Returns
-------
int
A stable 32-bit identifier for this node.
"""
try:
return self._stable_hash_value
except AttributeError:
content = repr(self.get_hashable()).encode("utf-8")
self._stable_hash_value = int(hashlib.md5(content).hexdigest()[:8], 16)
return self._stable_hash_value

def __hash__(self) -> int:
"""
Hash of an expression with caching.
Expand Down
47 changes: 36 additions & 11 deletions python/cudf_polars/cudf_polars/experimental/benchmarks/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,8 @@
if TYPE_CHECKING:
from collections.abc import Callable, Sequence

from cudf_polars.experimental.explain import SerializablePlan


try:
import structlog
Expand Down Expand Up @@ -237,6 +239,7 @@ class RunConfig:
default_factory=PackageVersions.collect
)
records: dict[int, list[Record]] = dataclasses.field(default_factory=dict)
plans: dict[int, SerializablePlan] = dataclasses.field(default_factory=dict)
dataset_path: Path
scale_factor: int | float
shuffle: Literal["rapidsmpf", "tasks"] | None = None
Expand Down Expand Up @@ -507,28 +510,37 @@ def print_query_plan(
args: argparse.Namespace,
run_config: RunConfig,
engine: None | pl.GPUEngine = None,
) -> None:
*,
print_plans: bool = True,
) -> tuple[str | None, str | None]:
"""Print the query plan."""
logical_plan = plan = None
if run_config.executor == "cpu":
if args.explain_logical:
print(f"\nQuery {q_id} - Logical plan\n")
print(q.explain())
logical_plan = q.explain()
if args.explain:
print(f"\nQuery {q_id} - Physical plan\n")
print(q.show_graph(engine="streaming", plan_stage="physical"))
plan = q.show_graph(engine="streaming", plan_stage="physical")
elif CUDF_POLARS_AVAILABLE:
assert isinstance(engine, pl.GPUEngine)
if args.explain_logical:
print(f"\nQuery {q_id} - Logical plan\n")
print(explain_query(q, engine, physical=False))
logical_plan = explain_query(q, engine, physical=False)
if args.explain and run_config.executor == "streaming":
print(f"\nQuery {q_id} - Physical plan\n")
print(explain_query(q, engine))
plan = explain_query(q, engine)
else:
raise RuntimeError(
"Cannot provide the logical or physical plan because cudf_polars is not installed."
)

if print_plans:
if logical_plan:
print(f"\nQuery {q_id} - Logical plan\n")
print(logical_plan)
if plan:
print(f"\nQuery {q_id} - Physical plan\n")
print(plan)

return logical_plan, plan


def initialize_dask_cluster(run_config: RunConfig, args: argparse.Namespace): # type: ignore[no-untyped-def]
"""
Expand Down Expand Up @@ -947,6 +959,12 @@ def parse_args(
help="Print an outline of the logical plan",
default=False,
)
parser.add_argument(
"--print-plans",
action=argparse.BooleanOptionalAction,
help="Print the query plans",
default=True,
)
parser.add_argument(
"--validate",
action=argparse.BooleanOptionalAction,
Expand Down Expand Up @@ -1052,6 +1070,7 @@ def run_polars(
run_config = dataclasses.replace(run_config, n_workers=actual_n_workers)

records: defaultdict[int, list[Record]] = defaultdict(list)
plans: dict[int, SerializablePlan] = {}
engine: pl.GPUEngine | None = None

if run_config.executor != "cpu":
Expand Down Expand Up @@ -1079,7 +1098,13 @@ def run_polars(
except AttributeError as err:
raise NotImplementedError(f"Query {q_id} not implemented.") from err

print_query_plan(q_id, q, args, run_config, engine)
print_query_plan(
q_id, q, args, run_config, engine, print_plans=args.print_plans
)
if (args.explain or args.explain_logical) and engine is not None:
from cudf_polars.experimental.explain import serialize_query

plans[q_id] = serialize_query(q, engine)

records[q_id] = []
for i in range(args.iterations):
Expand Down Expand Up @@ -1140,7 +1165,7 @@ def run_polars(
)
records[q_id].append(record)

run_config = dataclasses.replace(run_config, records=dict(records))
run_config = dataclasses.replace(run_config, records=dict(records), plans=plans)

# consolidate logs
if _HAS_STRUCTLOG and run_config.collect_traces:
Expand Down
Loading
Loading