[Data] Break "UDF time" down into input prep, function body and output build - #65807
Conversation
64b00c8 to
e7d29a7
Compare
There was a problem hiding this comment.
Code Review
This pull request decomposes the overall 'UDF time' metric in Ray Data into four distinct phases: input preparation, UDF body execution, output block building, and other fused stages. It introduces these granular metrics across execution stats, operator summaries, and block metadata, with an opt-in configuration (accurate_map_phase_timing) for row-based transforms to avoid overhead. The review feedback highlights a critical missing update to BlockExecStatsBuilder.build that would cause a runtime TypeError when constructing block stats, and suggests a formatting improvement to ensure core UDF phases are consistently displayed in verbose logs even when their sum is zero.
…t build "UDF time" covers a map stage's whole transform, not the user's function: the timed window spans turning input blocks into batches, calling the function, and assembling the output back into blocks. On a batch that carries Python objects the two ends can dominate -- a UDF that never touches an object column still reports the time spent unpickling it -- and there is no way to tell which part is which. Decompose it. `udf_time_s` keeps meaning the whole chain, so nothing that reads it today changes, and four new figures say where inside it the time went: input prep, the UDF body, output block build, and the bodies of non-UDF stages fused into the same chain. They sum back to the total, which a test asserts. `ds.stats()` renders them as a breakdown under the existing line, and all four are `metric_field`s on `OpRuntimeMetrics`, so they reach Prometheus per operator the same way the existing task metrics do. Timing a phase reuses the nesting subtraction already in `_UDFTimingIterator`: it measures inclusive time and subtracts what the timers below it added, so splitting one timer per stage into three per stage needs no new mechanism, only new wrap points inside `MapTransformFn.__call__`. Row-based transforms are not decomposed by default. Each wrapper costs a Python frame per item it yields, which is noise for a batch transform and roughly 1us per row across a three-stage chain for a row transform. Those keep the existing one-timer-per-stage arrangement -- same cost, same total, no breakdown -- and `DataContext.accurate_map_phase_timing` opts into the breakdown for them. Note the timers stay per stage even without decomposition rather than collapsing to one for the whole chain: `MapTransformFn.__call__` runs `_pre_process` eagerly, so with `batch_size="auto"` a later stage pulls data through earlier ones while the chain is still being built, and a single timer installed at the end would miss it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Three cases: - The components add up to `udf_time`. This is what makes the breakdown a decomposition rather than four loosely related numbers, and it is the property that would break first if a phase were double counted or missed. - Input prep and output block build are visible separately from the body. The UDF reads only the `id` column and never touches the column of Python objects beside it, so input prep dominates -- time the single figure used to attribute to user code. - Row transforms report only a total by default, and flipping `accurate_map_phase_timing` buys the breakdown without moving the total. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
…e line Three changes from review of the phase breakdown. The breakdown no longer appears in the default `ds.stats()`. It is rendered only under `verbose_stats_logs`, the same flag and the same treatment `extra_metrics` already gets, so the default output is byte-identical to before and `test_dataset_stats_basic` passes unchanged. The figures are still always on `OperatorStatsSummary` for anyone reading it programmatically -- gating display is not a reason to withhold data. "Other stages" is now "Built-in stages". The old name said what the line was not, which collided with its three siblings: they are all non-UDF work too. The new one says what it is -- the body of a stage Ray Data supplied, as against the function you passed. Reads, writes, downloads and file listing all land there, so the label generalises where an enumeration would not. The phase figures are `None` rather than zero when a chain measured only its total, so a consumer can tell "not measured" from "measured as zero", and the five metrics are `map_only`: an operator with no UDF transform chain has no UDF phases to report, and emitting zeros for it is noise. Tests: the phase timings are canonicalized to a stable placeholder like `obj_store_mem_used`. A trivial UDF spends sub-microsecond in each phase, so the figure rounds to zero or not depending on the run -- observed flipping between two consecutive runs. The measured-vs-not distinction that would have tested is asserted directly in `test_row_transform_phases_are_opt_in` instead. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
… non-UDF operators Four review findings on this PR, two of them real bugs. Non-UDF operators reported UDF time. The `decomposed` branch wrapped every stage in a timer, so a standalone read or projection reported its whole transform under "UDF time" -- an operator that runs nothing the user wrote. The parent branch and master both install a timer only `if _is_udf`, and report zero for those operators. `apply_transform` now skips timing entirely when no UDF is in the chain, so there is no timer to pay for either. Eagerly-consuming stage bodies went untimed. `generate_write_fn` has no `yield`: it calls `datasink.write(it1, ctx)`, which performs the whole upload before returning the iterable. That work happened inside `_apply_transform`, which was evaluated before `wrap_phase` had anything to wrap, so it landed in no phase and not in the total. A datasink sleeping 0.3s per block over two blocks reported 155us of built-in stage time against 600ms of I/O. `PhaseWrapFn` now takes a thunk instead of an iterable, so the hook owns the call and can time it. `UDFTimeScope.attribute_call` times that call with the same nesting subtraction as `__next__` -- an eager stage pulls its input through the timers already in the chain, so subtracting what they credited leaves its own work. The window runs once per stage per task, off the per-row path, so it applies on both the decomposed and total-only paths and the two totals stay in agreement. The same workload now reports 606ms, and UDF time accounts for 99.5% of remote wall time rather than 25.6%. The three phases every chain runs are now printed even at 0.0, so the lines visibly sum to the total above them. Previously a phase was dropped at zero, which both broke that arithmetic and made line presence vary run to run for a fast UDF -- a latent flaky test. Built-in stages is still dropped at zero, since most chains fuse nothing built-in. Docs: the operator-stats section described "UDF time" as excluding input prep and output block build. It includes them -- they are its components, and `Function body` is the line that excludes them. Restructured to nest the four phases under the total they decompose, documented `Function body` and `Built-in stages`, and noted that the breakdown is verbose-only and that row transforms need `accurate_map_phase_timing`. Added it to the Verbose stats list alongside `extra_metrics`, which it follows. Tests: `test_operator_without_a_udf_reports_no_udf_time` and `test_eagerly_consuming_stage_body_is_timed`, each confirmed to fail before being confirmed to pass. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
84cdde4 to
72281ba
Compare
iamjustinhsu
left a comment
There was a problem hiding this comment.
Nice, left some questions and some higher-level feedback since I think the code is getting a bit complicated
Two review nits from @iamjustinhsu, neither functional. The "UDF time" bullet spelled out that the figure covers "the work Ray Data does on either side of them to hand them batches and to turn what they return back into blocks", which is a mouthful for something the breakdown right below it already itemises. Shortened to say it covers the functions plus the work around them that feeds and collects them, and stated up front that the figure is per output block -- a second question the reviewer raised about the same paragraph. The comment above the breakdown table had grown to 17 lines for 12 lines of code, restating rationale that lives in the PR description and in the `accurate_map_phase_timing` docstring. Cut to the three things a reader of this function cannot get elsewhere: the lines sum to the total, they are verbose-only, and the trailing bool is whether to print at 0.0. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Answers a review comment on #65807 asking for an approach that is explicit about what gets timed, rather than one where the timing emerges from nested lazy iterators. `MapTransformFn.timed_steps()` returns the transform as a flat list of `TimedStep`s -- label, phase, stage index, and an `Iterable -> Iterable` body. `MapTransformer.get_timed_steps()` concatenates them and `TransformClock.chain()` builds the pipeline, so the answer to "what is getting timed?" is a list that can be printed and asserted on. A step measures only itself, with no coordination: def __next__(self): start = time.perf_counter() try: if self._iter is None: self._iter = iter(self._apply(self._upstream)) return next(self._iter) finally: self._totals[self._idx] += time.perf_counter() - start That total is inclusive of everything upstream. The chain is linear -- step k is only ever pulled by step k+1 -- so all of a step's time lies inside its consumer's windows, and subtracting neighbouring totals at drain recovers each step's own work. #65807 does the same arithmetic per item, threading a shared cursor through every `__next__`; doing it once over a list is easier to follow and measures 580 ns/item against 1016. Deferring `apply` to the first pull removes the eager special case. A write performs its whole upload while its iterable is built, which on #65807 needed its own timing window (`attribute_call`); here that work simply happens inside the same window as the pulls. `accurate_map_phase_timing` stops being a second code path and becomes a rewrite of the step list: fold a stage's three steps into one, whose `bucket` is None because it spans all three. A phase no step carries is then absent from the grouping, so "not measured" falls out of the data rather than being tracked by a flag. Accruing per step also means the per-phase and per-fused-stage breakdowns are two groupings of one array, with no second accumulator to keep in sync. `UDFTimeScope` becomes `TransformClock` and the `udf_time_scope=` keyword becomes `clock=`; the old name stopped describing an object that times every step, including the ones that are not user code. Removes `PhaseWrapFn`, `_no_phase_wrapping`, `wrap_phase`, `attribute_call` and `_UDFTimingIterator`. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
`from ray.data.datasource import Datasink`, added with the eager-stage test, sat above the `ray.data._internal` block. Ray's pre-commit runs ruff a second time with `--select I`, which is not in the default select set, so a plain `ruff check` passes while CI does not. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Answers a review comment on #65807 asking for an approach that is explicit about what gets timed, rather than one where the timing emerges from nested lazy iterators. `MapTransformFn.timed_steps()` returns the transform as a flat list of `TimedStep`s -- label, phase, stage index, and an `Iterable -> Iterable` body. `MapTransformer.get_timed_steps()` concatenates them and `TransformClock.chain()` builds the pipeline, so the answer to "what is getting timed?" is a list that can be printed and asserted on. A step measures only itself, with no coordination: def __next__(self): start = time.perf_counter() try: if self._iter is None: self._iter = iter(self._apply(self._upstream)) return next(self._iter) finally: self._totals[self._idx] += time.perf_counter() - start That total is inclusive of everything upstream. The chain is linear -- step k is only ever pulled by step k+1 -- so all of a step's time lies inside its consumer's windows, and subtracting neighbouring totals at drain recovers each step's own work. #65807 does the same arithmetic per item, threading a shared cursor through every `__next__`; doing it once over a list is easier to follow and measures 580 ns/item against 1016. Deferring `apply` to the first pull removes the eager special case. A write performs its whole upload while its iterable is built, which on #65807 needed its own timing window (`attribute_call`); here that work simply happens inside the same window as the pulls. `accurate_map_phase_timing` stops being a second code path and becomes a rewrite of the step list: fold a stage's three steps into one, whose `bucket` is None because it spans all three. A phase no step carries is then absent from the grouping, so "not measured" falls out of the data rather than being tracked by a flag. Accruing per step also means the per-phase and per-fused-stage breakdowns are two groupings of one array, with no second accumulator to keep in sync. `UDFTimeScope` becomes `TransformClock` and the `udf_time_scope=` keyword becomes `clock=`; the old name stopped describing an object that times every step, including the ones that are not user code. Removes `PhaseWrapFn`, `_no_phase_wrapping`, `wrap_phase`, `attribute_call` and `_UDFTimingIterator`. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Justin asked, on the verbose breakdown bullet, what UDF time means when a task has several input or output blocks. It is measured per output block: min, max and mean are across blocks, the total is the operator's. Say so in both bullets -- the one he commented on and the primary definition, which also now distinguishes UDF time from the task's total time. Signed-off-by: Marwan Sarieddine <marwan@anyscale.com> Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FfAW6hXHCMinYkJ3W7nfBa Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
…f steps Folds the RFC branch in. It answers the review comment asking for an approach that is explicit about what gets timed, rather than one where the timing emerges from nested lazy iterators, and it replaces the per-item subtraction cursor this branch had. `MapTransformFn.timed_steps()` returns the transform as a flat list of `TimedStep`s -- label, bucket, stage index, and an `Iterable -> Iterable` body. `MapTransformer.get_timed_steps()` concatenates them and `TransformClock.chain()` builds the pipeline, so "what is getting timed?" is answered by a list you can print and assert on. Each step measures only itself, with no cross-step coordination: a step's total is inclusive of everything upstream, the chain is linear, so subtracting neighbouring totals once at drain recovers each step's own work. That is 580 ns/item against the 1016 ns/item the shared cursor cost. Deferring `apply` to the first pull also removes the eager special case: a write performs its whole upload while its iterable is built, which previously needed its own timing window (`attribute_call`); here that work lands inside the same window as the pulls. Signed-off-by: Marwan Sarieddine <marwan@anyscale.com>
…nding Cursor Bugbot caught this on the merge of #65855. `TransformClock.chain` hands every `_TimedStep` a reference to `self.inclusive`, and `drain` replaced that list with a new one -- so after the first drain the steps kept adding to a list the clock no longer read. `_map_task` drains once per output block, so a task's first block reported its time and every block after it reported zero. Zeroing the list in place keeps the steps and the clock pointing at the same object. The per-drain window is still self-contained: both a step and its upstream neighbour start from zero at the same instant, so the adjacent-difference arithmetic is unchanged. Unit-level, one stage sleeping 0.05s with one output block per batch: before block 0: 0.0573s block 1: 0.0000s block 2: 0.0000s after block 0: 0.0826s block 1: 0.0537s block 2: 0.0548s `test_every_output_block_is_timed` covers it, and fails on the old line. Signed-off-by: Marwan Sarieddine <marwan@anyscale.com> Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FfAW6hXHCMinYkJ3W7nfBa Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Vale is what fails microcheck's lint job, and all four of its errors were
in lines this PR added:
396:93 Google.EmDash Don't put a space before or after a dash.
396:94 Google.EnDash Use an em dash ('---') instead of '--'.
397:50 Google.Contractions Use 'it's' instead of 'It is'.
464:73 Google.WordList Use 'preceding' instead of 'above'.
Rather than switch to `---`, drop the dash: this file has no em dashes in
prose, so the only one would have been mine. The "described above"
cross-reference goes too -- the phases are named in the same sentence, so
it added nothing. Both bullets are active voice now, which also clears
three of Vale's suggestions on the same lines.
`pre-commit run vale --all-files` passes (it needs `rst2html` from
docutils on PATH, which is why this went unnoticed locally).
Signed-off-by: Marwan Sarieddine <marwan@anyscale.com>
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01FfAW6hXHCMinYkJ3W7nfBa
Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Justin asked, of a `Read -> select_columns -> map1 -> map2` pipeline, whether Function body is just map1 + map2, whether Built-in stages is Read + select_columns, and where the batching in between lands. I answered on the thread and said the doc should carry it; this is that. The new "Reading the UDF time breakdown" section runs his exact pipeline and reads the output line by line: Function body holds the two functions alone (616 ms against 0.6 s slept), Built-in stages holds ReadRange and Project, and Input prep and Output block build cover all four stages rather than only the two the caller wrote. Numbers are measured, not illustrative. `MapTransformPhaseTimes` also now says outright that `total_s` is not the function-body time and that `udf_body_s` is -- the ambiguity behind his rename request, stated where he raised it. Signed-off-by: Marwan Sarieddine <marwan@anyscale.com> Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FfAW6hXHCMinYkJ3W7nfBa Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Justin's point: the name reads as the function-body time, and it isn't. The figure covers a map stage's whole transform -- batch formation, the stage bodies, and output block building -- for every stage of a (possibly fused) chain. `block_transform_time_s` says what it measures: the time to transform one output block. `ds.stats()` prints "Block transform time". This is a user-visible rename, so it lands on its own: - `BlockExecStats.udf_time_s` -> `block_transform_time_s` - `OperatorStatsSummary.udf_time` -> `block_transform_time` - `OpRuntimeMetrics.udf_time_s` -> `block_transform_time_s`, so the Prometheus metric is `ray_data_block_transform_time_s` - the `* UDF time:` line in `ds.stats()` -> `* Block transform time:` The metric is the reason to do it now rather than later: this PR introduces it, so nothing is scraping `ray_data_udf_time_s` yet. The other three names predate the PR, and renaming them changes a printed label, a summary attribute and an `extra_metrics` key. `block.py` and `MapTransformPhaseTimes` keep a note of the old name, so the rename is discoverable from the code rather than only from git. `udf_body_time_s` keeps its name here: it really is the UDF bodies. The next commit widens it and renames it accordingly. Signed-off-by: Marwan Sarieddine <marwan@anyscale.com> Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FfAW6hXHCMinYkJ3W7nfBa Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Justin's side note on the restructure thread: built-in stages should be
treated the same as UDFs, since a read or a write is a function like any
other and doesn't need its own breakdown.
`MapTransformPhase.OTHER` is gone, so a stage body lands in
`FUNCTION_BODY` whoever wrote it. `udf_body_*` becomes `function_body_*`
throughout, because the figure no longer covers only UDF bodies, and
`other_stage_time_s` disappears from `BlockExecStats`,
`OpRuntimeMetrics` and `OperatorStatsSummary`. The breakdown is three
lines rather than four.
That drops the "print at 0.0?" flag the rendering carried: it existed
only so `Built-in stages` could hide itself on a chain with no built-in
stage fused in. All three remaining lines always print, so `breakdown`
is a list of pairs and the loop is one condition shorter.
The same pipeline as the doc's worked example, before and after:
before after
* Block transform time: 697.95ms total * Block transform time: 648.79ms total
* Input prep: 3.19ms * Input prep: 1.69ms
* Function body: 616.17ms * Function body: 645.66ms
* Output block build: 2.61ms * Output block build: 1.44ms
* Built-in stages: 75.99ms
What this costs, stated plainly: `Function body` now mixes ReadRange and
Project in with map1 and map2, so "my code or Ray Data's?" is no longer
answerable from the breakdown. `TimedStep.stage_idx` already carries what
answers it properly, per stage and by name, which is #65810.
Also renames a `udf_stats` local the previous commit missed.
Signed-off-by: Marwan Sarieddine <marwan@anyscale.com>
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01FfAW6hXHCMinYkJ3W7nfBa
Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes using default effort and found 1 potential issue.
❌ Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, have a team admin enable autofix in the Cursor dashboard.
Reviewed by Cursor Bugbot for commit c3e49a8. Configure here.
…ering
`apply_transform` used to skip timing when no stage in the chain was a
UDF, so a standalone read or write reported nothing. That made sense when
the figure was called "UDF time"; under `block_transform_time_s` it reads
as a gap, and a read's body is a function like a UDF's. The gate is gone,
along with `MapTransformFn`'s `is_udf` parameter, which nothing else read
once `FUNCTION_BODY` stopped depending on it.
Removing the gate exposed a real ordering bug rather than just widening
coverage, and this is the part worth reviewing. The untimed path applied
every stage eagerly:
for step in steps:
data = step.apply(data)
`TransformClock.chain` instead defers each `apply` to the first pull,
which is what lets it measure a datasink that uploads inside `apply`. The
checkpoint write chain is prepare -> write -> commit, and
`commit_checkpoints` read `ctx.kwargs` **before** consuming its input, so
with the applies deferred it ran first, found no pending checkpoints, and
committed nothing. The pending checkpoint was written and never
committed, so recovery treated the run's data files as orphaned:
4 failed, 56 passed test_checkpoint.py
test_partial_failure_no_duplicates
test_partial_failure_no_duplicates_partitioned
test_2pc_fail_retry_cleans_pending_checkpoints
test_checkpoint_restore_after_full_execution
`commit_checkpoints` now drains upstream before reading `ctx`, which is
what "AFTER the data write succeeds" means -- the write is upstream, so
consuming its output is what waits for it. That is already how
`generate_collect_write_stats_fn` does it on the non-checkpoint write
path, which is why that one was unaffected. Back to 60 passed.
`test_operator_without_a_udf_reports_no_block_transform_time` becomes
`test_operator_without_a_udf_is_still_timed`, and asserts the phases are
measured, that they sum to the total, and that the total stays inside the
task's wall time.
Signed-off-by: Marwan Sarieddine <marwan@anyscale.com>
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01FfAW6hXHCMinYkJ3W7nfBa
Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Cursor Bugbot: the three phase metrics were summed with `or 0`, so a
default `map` or `filter` -- which is timed as a whole rather than phase
by phase -- exported three zeros beside a non-zero
`block_transform_time_s`. A dashboard couldn't tell that from "these
phases took no time", and the zeros don't add up to the total.
`BlockExecStats` and `OperatorStatsSummary` already keep the distinction;
only `OpRuntimeMetrics` flattened it. The three are `Optional[float]`
defaulting to `None` now, and accumulation skips a `None` rather than
adding zero for it, so they stay absent until a block reports a phase.
`Optional` metrics are already routine here -- `average_num_outputs_per_task`
and friends return `None` before any task finishes, through the same
export path.
Measured, same three pipelines:
row map (default) total=0.0024 prep=None body=None build=None
map_batches total=0.00052 prep=0.000168 body=0.000109 build=0.000239
row map (flag on) total=0.00054 prep=0.000082 body=0.000110 build=0.000350
The two measured rows sum to their totals; the first no longer claims to.
`test_phase_metrics_stay_none_when_unmeasured` covers it and fails on the
old `or 0`, with `test_phase_metrics_accumulate_when_measured` guarding
the other direction. The stats repr tests canonicalize `None` alongside
the numbers, as they already did to keep sub-microsecond figures stable.
Also documents the four metrics in the Task metrics table, which this PR
had added to Prometheus without listing.
Signed-off-by: Marwan Sarieddine <marwan@anyscale.com>
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01FfAW6hXHCMinYkJ3W7nfBa
Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
`TransformClock.chain` runs nothing until the result is pulled, which is what puts an eagerly consuming stage's work inside a timed window. The requirement that places on a stage was implicit, and the checkpoint 2PC break in c89db16 was the first thing to violate it: a stage that depends on an upstream stage's side effects has to consume upstream to get them. Stating it at the mechanism rather than only at the one victim, and naming both transforms that drain first so the next author has an example. Signed-off-by: Marwan Sarieddine <marwan@anyscale.com> Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FfAW6hXHCMinYkJ3W7nfBa Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
|
@iamjustinhsu this has moved a lot since your review — worth re-reading the description rather than the old threads.
Four things need your call; they're listed under What needs review in the description. |
Six lines to say one thing. The reader at that line only needs to know why `list(blocks)` isn't dead code, given the function returns `blocks` anyway: pulling upstream is what makes the prepare stage run. The rest of what the old comment said is on `TransformClock.chain`, where it applies to every stage rather than just this one. Signed-off-by: Marwan Sarieddine <marwan@anyscale.com> Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FfAW6hXHCMinYkJ3W7nfBa Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
dstrodtman
left a comment
There was a problem hiding this comment.
Docs-team style pass from Douglas Strodtman (Anyscale docs). Claude Code assisted; I read every comment below and stand behind each one.
Scope: prose style, grammar, and docs conventions only. I'm deliberately staying out of the naming (UDF time → Block transform time) and the phase-model questions in @iamjustinhsu's open threads — those are the Data team's call, and they look like they're still in motion. The new prose reads clearly; the notes below are all mechanical.
Four small fixes inline, all serial-comma and filler-word nits against the Google developer style guide (Ray's stated fallback for anything the Ray writing-style guide doesn't cover). Nothing here blocks — take or leave each individually. No technical review implied, and no approval: a Data maintainer is already reviewing.
Co-authored-by: Douglas Strodtman <douglas@anyscale.com> Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Co-authored-by: Douglas Strodtman <douglas@anyscale.com> Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Co-authored-by: Douglas Strodtman <douglas@anyscale.com> Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Co-authored-by: Douglas Strodtman <douglas@anyscale.com> Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
iamjustinhsu
left a comment
There was a problem hiding this comment.
Nice, I think once we add the data_dashboard panels it should be good, or a follow up would be fine too
| description="Time spent serializing blocks produced.", | ||
| metrics_group=MetricsGroup.TASKS, | ||
| ) | ||
| block_transform_time_s: float = metric_field( |
There was a problem hiding this comment.
I think we should also update the data_dashboard_panels.py so that we can get accurate metrics as well. The metric name appends the ray_data_ prefix, so in this case you have created ray_data_block_transform_time_s, etc... which can be queried with promQL
There was a problem hiding this comment.
Added in 175de2f. Two panels in the Outputs row next to Block Generation Time: Block Transform Time, and Block Transform Time by Phase stacking the three phases. Both divide by num_task_outputs_generated like the neighbouring panel, since these are per-output-block sums.
One thing to flag: the phase panel is empty for a row transform that did not opt into accurate_map_phase_timing. Those gauges stay None and Gauge.set skips None, so nothing is recorded rather than three misleading zeros. The panel description says so.
| # Upstream is lazy: nothing runs until its output is pulled. Pulling it | ||
| # here is what makes the prepare stage run and leave the pending | ||
| # checkpoints on `ctx` for the loop below to commit. | ||
| blocks = list(blocks) |
There was a problem hiding this comment.
Can you elaborate more on why this is being eagerly pulled? I don't quite understand the 2nd sentence
There was a problem hiding this comment.
Rewritten in ff53e7e:
# Each stage runs on the first pull from it, so nothing upstream has
# run yet. `prepare_checkpoint` is what leaves the pending checkpoints
# on `ctx`. Drain first, or there is nothing here to commit.Additionally, TransformClock.chain doc string now concretely explains how the chain is built and executed.
| batches = self._pre_process(blocks) | ||
| results = self._apply_transform(ctx, batches, report_custom_op_stats) | ||
| return self._post_process(results) | ||
| """Apply this transform's steps in order, untimed. |
There was a problem hiding this comment.
I think this docstring is slightly misleading: this function does time, but the timing is delegated to the TimeStep classes
There was a problem hiding this comment.
Fixed by renaming rather than rewording, in cbf57bb. A Step is inert — a body plus the metadata the summary groups by — so timed_steps() was the misleading part. It is Step / MapTransformFn.steps / MapTransformer.get_steps now, and _TimedStep stays as the wrapper that does the timing.
| # this list, so rebinding it here would leave them adding to a list | ||
| # this clock no longer reads -- and `_map_task` drains after every | ||
| # output block, so only a task's first block would report any time. | ||
| self.inclusive[:] = [0.0] * len(self.inclusive) |
There was a problem hiding this comment.
I think it makes more sense to avoid _TimeStep accepting a mutable reference to inclusive. We can probably just init to 0s in _TimedStep, remove this line, and then in _self_times, we can call some _TimedStep.drain() to get the inclusive times.
There was a problem hiding this comment.
Done in 2ca8f2e. _TimedStep holds its own total and hands it over through drain(); TransformClock collects them in order before subtracting neighbours, which it did anyway. The shared list and the in-place reset are both gone, along with the comment that was guarding them.
A `Step` is inert -- a body plus the metadata the summary groups by -- so calling it `TimedStep` made `__call__`, which composes steps with no clock anywhere, read as if it timed them. `_TimedStep` is the wrapper that does the timing. Renames the dataclass, `MapTransformFn.steps` and `MapTransformer.get_steps` to match. Also drops the two comments that only parse for a reader who saw the old `udf_time_s` name, and says what the commit stage's drain prevents: reading `ctx` before draining finds no pending checkpoints, commits nothing, and silently skips the second phase of the 2PC. Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
`_TimedStep` took a shared list and an index into it, so the clock had to reset that list in place -- rebinding it would have left every step adding to a list the clock no longer read, and `_map_task` drains after every output block, so only a task's first block would have reported any time. That was a real bug once, caught in review, and the fix needed a comment to stay fixed. Each step now holds its own total and hands it over through `drain()`. The clock collects them in order before subtracting neighbours, which it has to do anyway, so nothing about the arithmetic changes -- but there is no longer a shared object to reset, and the class of bug stops existing. Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
The new metrics export as `ray_data_block_transform_time_s` and the three phase gauges, so they were queryable but not plotted. Adds two panels to the Outputs row, beside Block Generation Time, which they sit next to in shape: both are per-output-block figures, so both divide by `num_task_outputs_generated` rather than plotting the operator's cumulative sum. The phase panel stacks its three targets, since they are a decomposition of the total rather than three independent series. It is empty for a row transform that did not opt into `accurate_map_phase_timing`: those phase gauges stay `None` and `Gauge.set` skips `None`, so nothing is recorded rather than three misleading zeros. The panel description says so, since an empty panel otherwise reads as a broken one. Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Three short steps instead of one long sentence: stages are lazy, prepare is what fills `ctx`, so drain before reading it. `TransformClock.chain` already names this function when it states the contract, so the pointer back the other way is not needed. Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
Two readers get caught by the same gap. On the dashboard, Block Generation Time and Block Transform Time plot the same blocks and never agree: generation is wall clock over the block's whole window, transform is only the chain inside it. The panel now says which is larger and why, and that shuffles and aggregations report neither of the phase figures. In the code, four callers pass `clock=TransformClock()` and never drain it, which reads as an oversight rather than a choice. Those tasks build their own `BlockExecStats` and have nowhere to put a transform time, so they discard it; `apply_transform` now says so where the parameter is documented. Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
The docstring described the pull order in prose, which reads as abstract right where a reader needs to be concrete. It now traces a checkpointed write simplified to prepare -> commit: what the loop builds, then a line-by-line run of that chain, ending where `commit` reads the pending checkpoints that `prepare` left on the `TaskContext` during its drain. `generate_collect_write_stats_fn` depends on the same ordering and had nothing saying so -- its drain was a list comprehension that happens to run before the `ctx` read below it. The docstring used to carry that fact from a distance; a comment at the drain itself carries it better, so the note moves there and the docstring keeps to its one example. Also drops the note about eagerly consuming stages -- `_TimedStep.__next__` already explains why `apply` is deferred, and this docstring only needs the consequence. Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>
) > **Stacked on #65807.** The base branch is `marwan/data-per-stage-map-timing`, so this diff shows only the change on top of it. It will be retargeted to `master` as the stack lands. ## Why are these changes needed? Ray Data fuses adjacent operators, so one operator's **Block transform time** can cover several of your functions. #65807 splits that time by *phase* — input prep, function body, output block build — which tells you what kind of work is slow. It still cannot tell you **which** of two fused functions the body time went to, and that is the question you actually have when a `sort` step and a `score` step share an operator: ``` Operator 1 ReadRange->MapBatches(sort)->MapBatches(score): 4 tasks executed * Block transform time: 227.87ms mean, 911.47ms total * Input prep: 4.24ms total * Function body: 905.06ms total <- which of the two? * Output block build: 2.17ms total ``` Breaking the fusion to find out means inserting a `materialize()` between the steps, which changes the thing you are measuring. ## What this change does Adds a `Stage <n>` line per fused stage: ``` * Block transform time: 224.88ms min, 230.36ms max, 227.87ms mean, 911.47ms total * Input prep: 166.83us min, 2.61ms max, 1.06ms mean, 4.24ms total * Function body: 223.03ms min, 229.75ms max, 226.26ms mean, 905.06ms total * Output block build: 436.08us min, 693.87us max, 543.32us mean, 2.17ms total * Stage 0: 129.5us min, 357.71us max, 202.41us mean, 809.63us total * Stage 1: 22.75ms min, 26.6ms max, 24.49ms mean, 97.97ms total * Stage 2: 201.96ms min, 204.89ms max, 203.17ms mean, 812.69ms total ``` That run has stage 1 sleeping 0.02s and stage 2 sleeping 0.2s, and the split finds it. Stage 0 is the fused `ReadRange`. The two breakdowns are the same total cut two ways, and **both sum to it exactly**: 4.24 + 905.06 + 2.17 = 809.63us + 97.97 + 812.69 = 911.47ms. **The measurement was already there.** `TransformClock.drain` already turns each step's inclusive time into its own time and groups it by `Step.bucket`. Every step already carries a `stage_idx` as well, so grouping by that is a second dict in the loop that was already running — no new timers, and no change to the timing itself. The whole of `map_transformer.py` in this diff is that dict plus a flag and a docstring. `BlockExecStats` carries a per-stage tuple, and `OperatorStatsSummary` aggregates one `StatsSummary` per stage index. ### Off by default This puts one number per stage on every block's metadata, which is built and shipped per output block, so `DataContext.per_stage_map_timing` gates collection. An unfused operator reports a single entry equal to its total rather than nothing, so the line does not appear and disappear as fusion changes around it. **It is independent of `accurate_map_phase_timing`.** A row transform timed a stage at a time is still one step per stage carrying one stage index, so a row-based `map` can have the per-stage split without paying for the per-phase one: ``` * Block transform time: 274.12us min, 1.05ms max, 664.33us mean, 1.33ms total * Stage 0: 80.17us min, 118.25us max, 99.21us mean, 198.42us total * Stage 1: 117.37us min, 834.79us max, 476.08us mean, 952.17us total * Stage 2: 76.58us min, 101.5us max, 89.04us mean, 178.08us total ``` Rendering is gated a second time behind `verbose_stats_logs`, matching the phase breakdown, so the default `ds.stats()` output is unchanged even with the flag on. The figures are on `OperatorStatsSummary` either way, so `Dataset.get_stats_summary()` exposes them with no flag. ## Two limits worth knowing **Stages are identified by ordinal, not name.** `MapTransformer` holds a bare `List[MapTransformFn]`; the names live on the physical operator as `_logical_operators`, one layer up. So stages are numbered in chain order, matching the order they appear in the fused operator name (`ReadRange->MapBatches(sort)->MapBatches(score)` → stage 0 is the read, stage 1 is `sort`, stage 2 is `score`). Carrying real names through would mean plumbing them into the transformer, which seems worth doing separately if reviewers want it. **No Prometheus metric or dashboard panel.** The other figures on #65807 are fixed fields, so they map onto named gauges. A per-stage list is variable length and its meaning depends on which operators fused, so there is no stable metric name to give it. `ds.stats()` and `get_stats_summary()` are the interfaces here. Note also that where two `MapTransformFn`s are collapsed into a single one, one stage can correspond to two logical operators. ## Related issue number Follow-up to #65807. ## Checks - [x] I've signed off every commit (`git commit -s`). - [x] I've run `scripts/format.sh` to lint the changes in this PR. - [x] I've included any doc changes needed for https://docs.ray.io/en/master/. - [x] I've added any new APIs to the API Reference. For example, if I added a method in Tune, I've added it in `doc/source/tune/api/` under the corresponding `.rst` file. - [x] I've made sure the tests are passing. Note that there might be a few flaky tests, see the delay instructions if they timeout. - Testing Strategy - [x] Unit tests - [ ] Release tests - [ ] This PR is not tested :( ### AI assistance AI assistance was used for this change. It is not a duplicate: no open PR covers per-stage attribution inside a fused map operator — the only adjacent work is #65807, which this is stacked on and which splits by phase rather than by stage. Tests run locally: ``` pytest python/ray/data/tests/test_stats.py pytest python/ray/data/tests/test_map_transformer.py pytest python/ray/data/tests/test_op_runtime_metrics.py ``` --------- Signed-off-by: Marwan Sarieddine <sarieddine.marwan@gmail.com>

Why are these changes needed?
ds.stats()reports one figure per operator for a map stage's transform and calls it "UDF time". It isn't the user's function: it spans turning input blocks into batches or rows, calling the function, and assembling the output back into blocks. So a slow operator tells you nothing about whether to optimise your code or your data layout.That bites hardest when rows carry Python objects, because both ends of the window get expensive. A UDF that never reads an object column is still charged for unpickling it, and nothing in the output says so.
What this change does
Renames the figure to say what it measures, and breaks it down.
ds.stats()prints* Block transform time:where it printed* UDF time:, and underverbose_stats_logssplits it into three lines that sum to it:The total and all three phases are
metric_fields onOpRuntimeMetrics, so they reach Prometheus per operator taggeddatasetandoperator. The three phases stayNonerather than zero when a chain is timed only as a whole, so a dashboard can tell "not measured" from "took no time".Default output is unchanged: the breakdown renders only under
verbose_stats_logs, the same treatmentextra_metricsgets, and the figures are always onget_stats_summary()regardless. Row transforms (map,filter) report only a total by default, since a timer per row costs roughly 0.6µs across a three-stage chain;DataContext.accurate_map_phase_timingopts them in.The docs gain a worked example that runs a fused
ReadRange->Project->MapBatches(map1)->MapBatches(map2)and reads the three figures off its output, because "which stages land in which bucket?" was the first thing review asked.Under the hood
map_transformer.pyis rewritten by #65855, merged in here: a task's transform is now a flat list ofTimedSteps you can print and assert on, each timing itself with no coordination. Subtracting neighbouring totals once at drain recovers each step's own work, at 580 ns/item against the 1016 the per-item cursor cost.Timing is no longer gated on the chain containing a UDF, so reads and writes are measured too and
MapTransformFn.is_udfis gone, along with the separateBuilt-in stagesbucket. Removing the gate surfaced an ordering bug in the checkpoint 2PC chain, fixed in c89db16.What needs review
ray_data_block_transform_time_sis introduced here, so it is free to change now and expensive after a release. @iamjustinhsu suggestedtotal_task_per_block_s.Built-in stagesintoFunction bodycosts the "my code or Ray Data's?" split until per-stage attribution lands in [Data] Split block transform time per fused stage, behind a flag #65810. @iamjustinhsu asked for it; worth confirming the interim is fine.ds.stats()output for every read operator.Related issue number
None. Closed #55052 raised adjacent complaints about operator-level metric granularity and was addressed by
ds.explain(); it did not cover this.Checks
git commit -s).scripts/format.shto lint the changes in this PR.DataContext.accurate_map_phase_timingis documented in theDataContextdocstring; no new public API)Testing
Eight tests added:
test_block_transform_time_phases_sum_to_the_totaltest_block_transform_time_phases_separate_object_serdetest_row_transform_phases_are_opt_intest_eagerly_consuming_stage_body_is_timedFunction bodyinstead of being reported as microsecondstest_operator_without_a_udf_is_still_timedtest_map_transformer.py::test_every_output_block_is_timedtest_op_runtime_metrics.py::test_phase_metrics_stay_none_when_unmeasuredtest_op_runtime_metrics.py::test_phase_metrics_accumulate_when_measuredNonedefault doesn't swallow a measured phaseExisting tests that assert on the rendered
ds.stats()string keep their non-verbose expectations unchanged, which is what gating the breakdown buys; the verbose ones gained three lines. The newOpRuntimeMetricskeys are canonicalized to a placeholder in those expectations, asobj_store_mem_usedalready is, because a trivial UDF's phases round to zero or not depending on the run.Run locally against a source build, with
PYTHONPATH=python/ray/data/tests:test_stats.pytest_map.pytest_checkpoint.pytest_consumption.pytest_map_transformer.py,unit/test_auto_batch_size.py,test_op_runtime_metrics.py,test_context.pydatasource/{test_datasink,test_file_datasink,test_parquet}.pyEvery failure above was A/B'd against master and traces to something absent from my machine rather than to this change: the dashboard API server isn't running, pandas 3.0 breaks an untouched conversion helper, and the rest are s3 fixtures.
pre-commit run --all-filespasses.AI assistance
AI assistance (Claude Code) was used to investigate the attribution gap, write the change and the tests, and draft this description. Per
AGENTS.md, the author has reviewed every changed line and the tests reported above were run locally.Duplicate check
gh pr list --repo ray-project/ray --state openwas searched forudf time phase breakdown input prepandmap_transformer phase timing; neither returns anything. The open PRs touching nearby code — #61379 (caching the serialized map transformer), #64090 (per-actor placement groups), #65663 (exposing backpressure policy in diagnostics) — do not touch this timing.