Skip to content

[Data] Break "UDF time" down into input prep, function body and output build - #65807

Merged
goutamvenkat-anyscale merged 29 commits into
masterfrom
marwan/data-per-stage-map-timing
Sep 15, 2026
Merged

goutamvenkat-anyscale merged 29 commits into
masterfrom
marwan/data-per-stage-map-timing

Conversation

@marwan116

@marwan116 marwan116 commented Aug 31, 2026 •

Copy link
Copy Markdown
Contributor

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 under verbose_stats_logs splits it into three lines that sum to it:

* Block transform time: 4.02s min, 4.11s max, 4.06s mean, 16.2s total
	* Input prep: ...           forming batches or rows, where object unpickling lands
	* Function body: ...        the stage bodies, yours and Ray Data's alike
	* Output block build: ...   assembling output back into blocks

The total and all three phases are metric_fields on OpRuntimeMetrics, so they reach Prometheus per operator tagged dataset and operator. The three phases stay None rather 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 treatment extra_metrics gets, and the figures are always on get_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_timing opts 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.py is rewritten by #65855, merged in here: a task's transform is now a flat list of TimedSteps 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_udf is gone, along with the separate Built-in stages bucket. Removing the gate surfaced an ordering bug in the checkpoint 2PC chain, fixed in c89db16.

What needs review

  1. The name. ray_data_block_transform_time_s is introduced here, so it is free to change now and expensive after a release. @iamjustinhsu suggested total_task_per_block_s.
  2. The checkpoint fix wants a checkpoint owner. It belongs in this PR, since the bug is created by this PR's deferred apply and doesn't exist on master, but nobody else has looked at it.
  3. Folding Built-in stages into Function body costs 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.
  4. Reads and writes now report a figure where they reported nothing, which changes 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

  • I've signed off every commit (git commit -s).
  • I've run scripts/format.sh to lint the changes in this PR.
  • I've included any doc changes needed for https://docs.ray.io/en/master/.
  • I've added any new APIs to the API Reference. (DataContext.accurate_map_phase_timing is documented in the DataContext docstring; no new public API)
  • I've made sure the tests are passing.
  • Testing Strategy
    • Unit tests
    • Release tests
    • This PR is not tested :(

Testing

Eight tests added:

test asserts
test_block_transform_time_phases_sum_to_the_total the components add up, which is what makes this a decomposition rather than three loose numbers
test_block_transform_time_phases_separate_object_serde input prep dominates the body for a UDF that only counts rows beside an object column
test_row_transform_phases_are_opt_in row transforms report only a total by default, and the flag adds detail without moving it
test_eagerly_consuming_stage_body_is_timed a datasink's I/O reaches Function body instead of being reported as microseconds
test_operator_without_a_udf_is_still_timed a standalone read reports its phases, and they stay inside the task's wall time
test_map_transformer.py::test_every_output_block_is_timed every output block of a task is measured, not just the first
test_op_runtime_metrics.py::test_phase_metrics_stay_none_when_unmeasured a chain timed only as a whole exports no phases, rather than zeros beside a non-zero total
test_op_runtime_metrics.py::test_phase_metrics_accumulate_when_measured the None default doesn't swallow a measured phase

Existing 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 new OpRuntimeMetrics keys are canonicalized to a placeholder in those expectations, as obj_store_mem_used already 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:

file result
test_stats.py 105 passed, 2 skipped, 1 failed
test_map.py 252 passed, 1 failed
test_checkpoint.py 60 passed, 13 skipped, 4 errors
test_consumption.py 28 passed, 2 skipped, 3 failed
test_map_transformer.py, unit/test_auto_batch_size.py, test_op_runtime_metrics.py, test_context.py all passed
datasource/{test_datasink,test_file_datasink,test_parquet}.py failures identical to the merge base

Every 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-files passes.

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 open was searched for udf time phase breakdown input prep and map_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.

@marwan116
marwan116 force-pushed the marwan/data-per-stage-map-timing branch 2 times, most recently from 64b00c8 to e7d29a7 Compare August 31, 2026 19:21
@marwan116 marwan116 self-assigned this Aug 31, 2026
@marwan116
marwan116 marked this pull request as ready for review August 31, 2026 22:18
@marwan116
marwan116 requested review from a team as code owners August 31, 2026 22:18

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread python/ray/data/block.py Outdated
Comment thread python/ray/data/_internal/stats.py
Comment thread python/ray/data/_internal/execution/operators/map_transformer.py Outdated
Comment thread python/ray/data/_internal/execution/operators/map_transformer.py Outdated
@ray-gardener ray-gardener Bot added the data Ray Data-related issues label Sep 1, 2026
Base automatically changed from marwan/data-fix-fused-udf-time-double-count to master September 1, 2026 17:02
marwan116 and others added 4 commits September 1, 2026 10:02
…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>
@goutamvenkat-anyscale
goutamvenkat-anyscale force-pushed the marwan/data-per-stage-map-timing branch from 84cdde4 to 72281ba Compare September 1, 2026 17:02

@iamjustinhsu iamjustinhsu left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nice, left some questions and some higher-level feedback since I think the code is getting a bit complicated

Comment thread doc/source/data/monitoring-your-workload.rst Outdated
Comment thread doc/source/data/monitoring-your-workload.rst
Comment thread doc/source/data/monitoring-your-workload.rst Outdated
Comment thread python/ray/data/_internal/execution/operators/map_transformer.py Outdated
Comment thread python/ray/data/_internal/stats.py Outdated
Comment thread python/ray/data/_internal/execution/operators/map_transformer.py Outdated
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>
marwan116 added a commit that referenced this pull request Sep 2, 2026
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>
marwan116 and others added 4 commits September 1, 2026 19:46
`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>
Comment thread python/ray/data/_internal/execution/operators/map_transformer.py Outdated
marwan116 and others added 4 commits September 2, 2026 12:14
…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>

@cursor cursor Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Cursor Bugbot has reviewed your changes using default effort and found 1 potential issue.

Fix All in Cursor

❌ 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.

Comment thread python/ray/data/_internal/execution/interfaces/op_runtime_metrics.py Outdated
marwan116 and others added 3 commits September 2, 2026 14:59
…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>
@marwan116

marwan116 commented Sep 3, 2026 •

Copy link
Copy Markdown
Contributor Author

@iamjustinhsu this has moved a lot since your review — worth re-reading the description rather than the old threads.

  • [Data] RFC: build a map task's timing from a declared list of steps #65855 is merged in, so map_transformer.py builds on the prototype you suggested.
  • Both your requests are done:
    • udf_time is now block_transform_time_s
    • Built-in stages is now folded into Function body.
  • Removing the UDF gate surfaced a 2PC ordering bug in the checkpoint write chain. Fixed here, and it's the part I'd most like another pair of eyes on.

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 dstrodtman left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread doc/source/data/monitoring-your-workload.rst Outdated
Comment thread doc/source/data/monitoring-your-workload.rst Outdated
Comment thread doc/source/data/monitoring-your-workload.rst Outdated
Comment thread doc/source/data/monitoring-your-workload.rst Outdated
@marwan116 marwan116 added the go add ONLY when ready to merge, run all tests label Sep 8, 2026
marwan116 and others added 4 commits September 8, 2026 14:36
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 iamjustinhsu left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can you elaborate more on why this is being eagerly pulled? I don't quite understand the 2nd sentence

@marwan116 marwan116 Sep 11, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread python/ray/data/block.py Outdated
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.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think this docstring is slightly misleading: this function does time, but the timing is delegated to the TimeStep classes

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread python/ray/data/_internal/execution/operators/map_transformer.py Outdated
# 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)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread python/ray/data/_internal/execution/operators/map_transformer.py
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>
@goutamvenkat-anyscale
goutamvenkat-anyscale merged commit 5f80213 into master Sep 15, 2026
6 checks passed
@goutamvenkat-anyscale
goutamvenkat-anyscale deleted the marwan/data-per-stage-map-timing branch September 15, 2026 17:41
goutamvenkat-anyscale pushed a commit that referenced this pull request Sep 17, 2026
)

> **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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

data Ray Data-related issues go add ONLY when ready to merge, run all tests

4 participants