Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
82 commits
Select commit Hold shift + click to select a range
0efdf65
[WIP]: Add Quent Resource tracing to cudf-polars
TomAugspurger Jul 7, 2026
43ba6e5
many things
TomAugspurger Jul 8, 2026
a8b8eb7
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 13, 2026
43a2df9
is_io_node
TomAugspurger Jul 13, 2026
8745a5b
remove print
TomAugspurger Jul 13, 2026
0c5e505
change emission times
TomAugspurger Jul 13, 2026
323b6aa
Remove Actor-scoped tasks
TomAugspurger Jul 13, 2026
47993bb
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 13, 2026
bbc7966
Revert collective trace
TomAugspurger Jul 13, 2026
ba03d7c
simplify tests
TomAugspurger Jul 13, 2026
a20d3b9
docs
TomAugspurger Jul 13, 2026
3b7e6af
restore old tracing flow
TomAugspurger Jul 13, 2026
62c27ec
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 13, 2026
2fc4e74
imports
TomAugspurger Jul 13, 2026
f038b59
test skip
TomAugspurger Jul 14, 2026
681638e
tweak the test
TomAugspurger Jul 14, 2026
9652021
docs
TomAugspurger Jul 14, 2026
9821ef7
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 14, 2026
e02620c
Move network declaration
TomAugspurger Jul 14, 2026
c78bfc8
Engine-scoped network, channels
TomAugspurger Jul 14, 2026
8c4a72e
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 14, 2026
22bef5a
Fixes
TomAugspurger Jul 14, 2026
cd3fefc
docs
TomAugspurger Jul 14, 2026
23107f3
fixup
TomAugspurger Jul 14, 2026
dd17ab8
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 14, 2026
5ddc590
Refactor
TomAugspurger Jul 14, 2026
9442826
more refactor
TomAugspurger Jul 14, 2026
51729c4
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 14, 2026
5b61041
test size_bytes
TomAugspurger Jul 14, 2026
faafe6d
testing
TomAugspurger Jul 14, 2026
6186ae0
test coverage
TomAugspurger Jul 14, 2026
0b2a8f0
fixup
TomAugspurger Jul 15, 2026
545c33f
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 15, 2026
c6f38ac
Fixup
TomAugspurger Jul 16, 2026
740abaf
pass through
TomAugspurger Jul 16, 2026
615b082
Track Operator Statistics via Tracer
TomAugspurger Jul 17, 2026
e423554
Set tracer in IRExecutionContext
TomAugspurger Jul 17, 2026
a66bfe7
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 17, 2026
ba69bf0
Ignore new format
TomAugspurger Jul 17, 2026
af692d8
remove io_loading_at hck
TomAugspurger Jul 17, 2026
908f785
query_for docs
TomAugspurger Jul 17, 2026
83d5f9d
docstrings
TomAugspurger Jul 17, 2026
b098c60
Reuse get device memory
TomAugspurger Jul 17, 2026
42d0429
ir_context setting
TomAugspurger Jul 20, 2026
ed5cb9f
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 20, 2026
9206fa4
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 20, 2026
90edf42
task/chunk level stats
TomAugspurger Jul 20, 2026
9e72712
merge the conditions
TomAugspurger Jul 21, 2026
5fbd8f7
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 21, 2026
f535f70
instance name
TomAugspurger Jul 22, 2026
360b1a6
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 23, 2026
bc35212
fixups
TomAugspurger Jul 23, 2026
67d6198
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 24, 2026
f1f047e
cleanpu
TomAugspurger Jul 24, 2026
d933028
WorkerResources
TomAugspurger Jul 24, 2026
da24843
remove unused parameter
TomAugspurger Jul 24, 2026
9dc8d86
remove unused parameter
TomAugspurger Jul 24, 2026
bd7b29e
refactor
TomAugspurger Jul 24, 2026
0ebe378
cvoerage
TomAugspurger Jul 27, 2026
ccaf4aa
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 27, 2026
716109e
cleanup
TomAugspurger Jul 27, 2026
20483d4
revert
TomAugspurger Jul 27, 2026
76b12ae
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 28, 2026
90c624c
Update quent output
TomAugspurger Jul 28, 2026
3b8be32
Include type name in physical operator instance name
TomAugspurger Jul 28, 2026
b883dac
require tracing with --collect-traces
TomAugspurger Jul 28, 2026
c5a3f86
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Jul 29, 2026
e2fa3a1
Serialize all properties
TomAugspurger Jul 30, 2026
3f7f43c
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Aug 3, 2026
1e3f9dd
Merge branch 'main' into tom/quent-resources
TomAugspurger Aug 3, 2026
4d406df
coverage
TomAugspurger Aug 3, 2026
8b14654
Bump quent version
TomAugspurger Aug 4, 2026
90af2f9
Comm
TomAugspurger Aug 4, 2026
e3193eb
Reuse via closure
TomAugspurger Aug 4, 2026
7e60dc9
Cleanup exception handling
TomAugspurger Aug 4, 2026
ec42750
Updated size_bytes test
TomAugspurger Aug 4, 2026
15153e7
simplify int sizes
TomAugspurger Aug 4, 2026
49eb232
Revert "simplify int sizes"
TomAugspurger Aug 4, 2026
f82ceb9
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Aug 4, 2026
d0ca121
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Aug 5, 2026
6090c47
Fixed engine config serialization
TomAugspurger Aug 5, 2026
c1cdcce
Merge remote-tracking branch 'upstream/main' into tom/quent-resources
TomAugspurger Aug 7, 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
4 changes: 3 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -189,4 +189,6 @@ rmm_log.txt
python/cudf/cudf_pandas_tests/data/rmm_log.txt

# Quent traces
logs/*.ndjson
logs/**/*.ndjson
logs/**/*.qmi
logs/*.zip
4 changes: 4 additions & 0 deletions python/cudf_polars/cudf_polars/containers/dataframe.py
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,10 @@ def __init__(
self.table = plc.Table([c.obj for c in self.columns], num_rows=num_rows)
self.stream = stream

def _size_bytes(self) -> int:
"""Return the size of the dataframe in bytes."""
return sum(c.device_buffer_size() for c in self.table.columns())

def copy(self) -> Self:
"""Return a shallow copy of self."""
return type(self)(
Expand Down
15 changes: 15 additions & 0 deletions python/cudf_polars/cudf_polars/dsl/ir.py
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,8 @@

from cudf_polars.containers.dataframe import NamedColumn
from cudf_polars.dsl.utils.io import CachedParquetInfo
from cudf_polars.quent._context import QuentIRExecutionContext
from cudf_polars.streaming.actor_graph.tracing import ActorTracer
from cudf_polars.streaming.rank_aware_source import RankAwareSource
from cudf_polars.typing import CSECache, ClosedInterval, Schema, Slice as Zlice
from cudf_polars.utils.config import ParquetOptions
Expand Down Expand Up @@ -138,11 +140,17 @@ class IRExecutionContext:
A zero-argument callable that returns a CUDA stream.
query_id
Identifier for the query being executed.
quent_ir_execution_context
Optional Quent tracing context bound to a physical operator.
tracer
The actor tracer. Used to propagate statistics.
"""

py_executor: concurrent.futures.ThreadPoolExecutor | None = field(default=None)
get_cuda_stream: Callable[[], Stream] = field(default=get_cuda_stream)
query_id: uuid.UUID = field(default_factory=uuid.uuid4)
quent_ir_execution_context: QuentIRExecutionContext | None = None
tracer: ActorTracer | None = None

async def to_thread(
self, func: Callable[P, T], /, *args: P.args, **kwargs: P.kwargs
Expand Down Expand Up @@ -244,6 +252,9 @@ class IR(Node["IR"]):
schema: Schema
"""Mapping from column names to their data types."""

is_io_node: bool = False
"""Whether the node is an IO node."""

def get_hashable(self) -> Hashable:
"""
Hashable representation of node, treating schema dictionary.
Expand Down Expand Up @@ -697,6 +708,8 @@ class Scan(IR):
PARQUET_DEFAULT_CHUNK_SIZE: int = 0 # unlimited
PARQUET_DEFAULT_PASS_LIMIT: int = 16 * 1024**3 # 16GiB

is_io_node: bool = True

def __init__(
self,
schema: Schema,
Expand Down Expand Up @@ -1637,6 +1650,8 @@ class DataFrameScan(IR):
projection: tuple[str, ...] | None
"""List of columns to project out."""

is_io_node: bool = True

def __init__(
self,
schema: Schema,
Expand Down
53 changes: 48 additions & 5 deletions python/cudf_polars/cudf_polars/dsl/tracing.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@

import cudf_polars.containers
from cudf_polars.dsl import ir
from cudf_polars.dsl.ir import IRExecutionContext


class Scope(enum.StrEnum):
Expand Down Expand Up @@ -161,17 +162,19 @@ def log_do_evaluate(
if not LOG_TRACES:
return func
else: # pragma: no cover; requires CUDF_POLARS_LOG_TRACES=1
# do this just once
pynvml.nvmlInit()
maybe_handle = get_device_handle()
pid = _getpid()

@functools.wraps(func)
def wrapper(
cls: type[ir.IR],
*args: P.args,
**kwargs: P.kwargs,
) -> cudf_polars.containers.DataFrame:
# do this just once
pynvml.nvmlInit()
maybe_handle = get_device_handle()
pid = _getpid()
from cudf_polars.quent._types import Task

log = structlog.get_logger()

# By convention, all non-dataframe arguments (non-child) come first.
Expand All @@ -180,6 +183,23 @@ def wrapper(
list(args) + [v for k, v in kwargs.items() if k != "context"]
)[cls._n_non_child_args :] # type: ignore[assignment]

# And the kwonly 'context' argument has the IR execution context.
ir_execution_context: IRExecutionContext = kwargs["context"] # type: ignore[assignment]

if ir_execution_context.quent_ir_execution_context is not None:
quent_task = Task.from_ir(
cls, ir_execution_context.quent_ir_execution_context
)
ir_execution_context.quent_ir_execution_context.context._emit_task_begin_events(
cls,
quent_task,
ir_execution_context.quent_ir_execution_context,
input_frames_bytes=sum(frame._size_bytes() for frame in frames),
)
Comment thread
coderabbitai[bot] marked this conversation as resolved.

else:
quent_task = None

before_start = time.monotonic_ns()
before = make_snapshot(
cls, frames, phase="input", device_handle=maybe_handle, pid=pid
Expand All @@ -191,9 +211,27 @@ def wrapper(
# argument, followed by the method-specific arguments, and returns a DataFrame.

start = time.monotonic_ns()
result = func(cls, *args, **kwargs)
try:
result = func(cls, *args, **kwargs)
except Exception: # pragma: no cover;
result = None
raise
finally:
if (
quent_task is not None
and ir_execution_context.quent_ir_execution_context is not None
):
# TODO: This should emit some Chunk-level statistics (duration, rows, bytes, schema, etc.)
ir_execution_context.quent_ir_execution_context.context._emit_task_end_events(
cls,
quent_task,
ir_execution_context.quent_ir_execution_context,
result,
)
stop = time.monotonic_ns()

assert result is not None

after_start = time.monotonic_ns()
after = make_snapshot(
cls,
Expand All @@ -215,6 +253,11 @@ def wrapper(
)
log.info("Execute IR", **record)

if (tracer := ir_execution_context.tracer) is not None:
# ActorTracer.send updates row_count and chunk_count
tracer.input_bytes += sum(frame._size_bytes() for frame in frames)
tracer.output_bytes += result._size_bytes()

return result

return wrapper
Expand Down
38 changes: 33 additions & 5 deletions python/cudf_polars/cudf_polars/engine/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@
attach_cached_parquet_metadata,
prefetch_parquet_file_metadata_for_ir,
)
from cudf_polars.quent._plan import build_plan
from cudf_polars.quent._plan import build_plan, build_quent_operator_map
from cudf_polars.streaming.actor_graph.collectives import ReserveOpIDs
from cudf_polars.streaming.actor_graph.collectives.common import reserve_op_id
from cudf_polars.streaming.actor_graph.core import generate_network
Expand All @@ -54,8 +54,8 @@
from rapidsmpf.memory.buffer_resource import BufferResource
from rapidsmpf.streaming.core.context import Context

import cudf_polars.quent
import cudf_polars.quent._logging
import cudf_polars.quent._types
from cudf_polars.dsl.ir import IR
from cudf_polars.dsl.translate import Translator
from cudf_polars.quent._context import LocalQuentContext
Expand Down Expand Up @@ -461,6 +461,9 @@ def execute_ir_on_rank(
config_options: ConfigOptions[StreamingExecutor],
stats: StatsCollector,
collective_id_map: dict[IR, list[int]],
*,
quent_operator_map: dict[IR, cudf_polars.quent._types.Operator] | None = None,
local_quent_context: LocalQuentContext | None = None,
) -> tuple[DataFrame, list[ChannelMetadata]]:
"""
Execute a Polars IR query on a single rank's GPU.
Expand All @@ -487,6 +490,12 @@ def execute_ir_on_rank(
Statistics collector.
collective_id_map
Mapping from IR nodes to their pre-allocated collective operation IDs.
quent_operator_map
Mapping from IR nodes to their Quent operators, or ``None`` when tracing
is disabled.
local_quent_context
The local Quent context for this rank, or ``None`` when tracing is
disabled.

Returns
-------
Expand All @@ -507,6 +516,8 @@ def execute_ir_on_rank(
ir_context=ir_context,
collective_id_map=collective_id_map,
metadata_collector=metadata_collector,
quent_operator_map=quent_operator_map,
local_quent_context=local_quent_context,
)

try:
Expand Down Expand Up @@ -733,20 +744,34 @@ def evaluate_on_rank(
Collected channel metadata.
"""
stats = allgather_stats(comm, ctx.br(), ir, config_options, py_executor)
# ``get_stable_plan_id`` is a deterministic function of the IR
# structure, so every rank derives the same logical plan ID for a
# given query (only rank 0 emits the declaration, but physical plans
# on every rank reference it as their parent). It is *not* unique
# across collects, though: re-running an identical query would reuse
# the same plan ID under a different parent query. Namespacing by the
# per-collect ``query_id`` (which is identical across ranks but unique
# per collect) keeps the cross-rank agreement while making the plan ID
# unique per collect.
logical_plan_id = uuid.uuid5(query_id, str(ir.get_stable_plan_id()))

physical_op_by_id: dict[str, cudf_polars.quent._types.Operator] | None = None
quent_operator_map: dict[IR, cudf_polars.quent._types.Operator] | None = None

lowering, node_map = lower_ir_graph_with_node_map(
ir, config_options, stats, rank=comm.rank, nranks=comm.nranks
)
optimized = lowering.optimized
ir = lowering.lowered
partition_info = lowering.partition_info
# TODO: figure out if we emit anything about optimized.

if config_options.executor.quent_context is not None:
assert local_quent_context is not None
logical_plan_id = optimized.get_stable_plan_id()
plan, ops, ports, logical_op_by_id = build_plan(
optimized,
config_options,
query=local_quent_context.context.query,
query=local_quent_context.query,
plan_id=logical_plan_id,
worker=local_quent_context.worker,
instance_name="logical",
Expand All @@ -764,7 +789,7 @@ def evaluate_on_rank(
if config_options.executor.quent_context is not None:
assert local_quent_context is not None
physical_plan_id = uuid.uuid4()
local_quent_context.context._emit_physical_plan_events(
physical_op_by_id = local_quent_context.context._emit_physical_plan_events(
local_quent_context.logger,
ir,
config_options,
Expand All @@ -774,6 +799,7 @@ def evaluate_on_rank(
node_map=node_map,
logical_op_by_id=logical_op_by_id,
)
quent_operator_map = build_quent_operator_map(ir, physical_op_by_id)
ir_context = IRExecutionContext(
py_executor, get_cuda_stream=ctx.br().stream_pool.get_stream, query_id=query_id
)
Expand All @@ -796,6 +822,8 @@ def evaluate_on_rank(
config_options,
stats,
collective_id_map,
quent_operator_map=quent_operator_map,
local_quent_context=local_quent_context,
)


Expand Down
Loading
Loading