Skip to content

[Data] Suppress cluster scale-up when pipeline is backpressured by iteration - #66569

Open
Yonghui-Lee wants to merge 6 commits into
ray-project:masterfrom
Yonghui-Lee:issue-45331-autoscaler-iteration-backpressure
Open

Yonghui-Lee wants to merge 6 commits into
ray-project:masterfrom
Yonghui-Lee:issue-45331-autoscaler-iteration-backpressure

Conversation

@Yonghui-Lee

@Yonghui-Lee Yonghui-Lee commented Sep 29, 2026 •

Copy link
Copy Markdown

Description

When a streaming Ray Data pipeline is read by a slow consumer, for example a training loop reading iter_batches / iter_torch_batches / streaming_split, Ray Data keeps asking for more nodes even though the consumer is the bottleneck. The new nodes fill up with tasks that stall on the same consumer, so they cost money without adding throughput. This happens with both RateBasedClusterAutoscaler and DefaultClusterAutoscalerV2.

Root cause

Both autoscalers decide whether to scale up from RollingLogicalUtilizationGauge, a rolling average of global_usage / global_limits. V2 scales once any resource reaches 75%; the rate-based autoscaler scales unless all resources are below 50%. When the consumer is slow, the ops at the end of the pipeline are blocked on it, but their usage still counts:

  • Ops in task output backpressure keep the CPUs/GPUs of their paused tasks.
  • Their unread outputs (output queues plus the iterator's prefetch buffers) push object store utilization to about 100%. This also happens when the ops' tasks have finished and the ops sit idle with their outputs queued.
    More nodes raise the budget, so more tasks start and stall, which leads to another request. The rate-based autoscaler then asks for 2x the current resources, because its per-task output rates leave out the time a task waits for its outputs to be read.

Fix

When computing utilization for autoscaling decisions only, leave out the usage of ops that are blocked on a slow consumer:

  • An op is blocked if it's in task output backpressure, or if it has no running tasks and can't submit more even though its CPU/GPU/memory budget fits another task (e.g., because its finished outputs fill its object store budget).
  • A blocked op is left out only if every downstream op is left out too, i.e., the chain of blocked ops ends at the consumer. Downstream ineligible ops (e.g., Limit, OutputSplitter) are left out with their upstream op.

Why compute-bound pipelines still scale. An op only goes into output backpressure when something downstream isn't draining it. If that downstream thing is a compute-starved op, that op's own usage still counts. For example, in a CPU → GPU pipeline that is GPU-bound, only the CPU op is left out, so GPU utilization still drives scale-up. If the downstream thing is the external consumer, more nodes can't help, and nothing is requested.

Related issues

Fixes #45331

Additional information

Tests

  • test_resource_manager.py::test_global_usage_excluding_output_backpressure: covers no backpressure, slow consumer (terminal op and downstream Limit left out), and GPU-bound (downstream GPU op's usage kept).
  • test_autoscaler_iteration_backpressure.py (new, registered in BUILD.bazel):
    • gauge test with the switch on and off; exported metrics stay raw;
    • end to end: no scale-up requests while the consumer is paused;
    • end to end: a CPU-bound pipeline with a fast consumer still scales up.
  • With the kill switch off, the paused-consumer test fails and the CPU-bound test still passes, so the test does catch the bug.

AI assistance: AI assistance was used to analyze the issue and to write the code and tests. I reviewed every changed line and ran the tests above locally.

@Yonghui-Lee
Yonghui-Lee requested a review from a team as a code owner September 29, 2026 06:06

@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 prevents Ray Data from scaling up the cluster when the pipeline is backpressured by a slow consumer. It introduces a new method get_global_usage_excluding_output_backpressure in ResourceManager to exclude operators that are currently blocked by task output backpressure, along with their downstream ineligible operators, from the scaling resource usage calculation. This behavior is controlled by the environment variable RAY_DATA_AUTOSCALING_EXCLUDE_OUTPUT_BACKPRESSURED_USAGE. Additionally, corresponding unit tests have been added to verify the resource utilization gauge and autoscaling behavior under backpressure. There are no review comments, and I have no feedback to provide.

…eration

Signed-off-by: Li Yonghui <yonghuili@google.com>
@Yonghui-Lee
Yonghui-Lee force-pushed the issue-45331-autoscaler-iteration-backpressure branch from bc35ce8 to 5e5200a Compare September 29, 2026 06:12
@ray-gardener ray-gardener Bot added data Ray Data-related issues community-contribution Contributed by the community labels Sep 29, 2026
@iamjustinhsu

Copy link
Copy Markdown
Contributor

Hi @Yonghui-Lee, thanks for the PR! I read the description and took a look at the PR and have 2 questions:

  1. Wouldn't operators that are output backpressured be unblocked after the cluster scales up, suggesting that we should include those operators? Basically I'm wondering why are operators still in output backpressure after cluster autoscaling
  2. Would u be able to try the RateBasedClusterAutoscaler? It might be more conservative when scaling by computing the output production rate of each operator

…caler-iteration-backpressure

Signed-off-by: Li Yonghui <yonghuili@google.com>
@Yonghui-Lee

Copy link
Copy Markdown
Author

@iamjustinhsu Thanks for taking a look.

1. Why are ops still in output backpressure after scaling up?

It depends on what they're waiting on. In GPU training (the case in #45331), the last ops wait on the trainer, which reads one batch per worker per step at a pace set by its GPUs. More nodes don't make the trainer faster; they only give the pipeline a bigger budget, which it fills within a second.

2. RateBasedClusterAutoscaler

Thanks for the suggestion! I tried it on a 4-CPU cluster with the consumer paused after the first batch (like a trainer saving a checkpoint):

Autoscaler Without this PR With this PR
RateBasedClusterAutoscaler 14/14 decisions scale up, each asking for 8 CPUs 0/14
DefaultClusterAutoscalerV2 14/14 0/14

RateBasedClusterAutoscaler decides whether to scale with the same RollingLogicalUtilizationGauge as V2, which reads near 100%. And it sizes requests by per-task output rates, which leave out the time a task waits for its outputs to be read, so stalled tasks still look fast and it asks for 2x the current CPUs. Since the fix is in the usage that feeds the shared gauge, it covers both autoscalers.

Signed-off-by: Li Yonghui <yonghuili@google.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.

Stale Bugbot comment from a previous run.

Comment thread python/ray/data/_internal/execution/resource_manager.py Outdated
Signed-off-by: Li Yonghui <yonghuili@google.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.

Left some feedback, lmk what u think

An op is excluded if it's blocked on downstream and all its downstream ops
are excluded too, i.e., the chain of blocked ops ends at the consumer.
"""
excluded_ops = set()

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.

We can remove the excluded_ops set, and just iterate in reverse order (skipping ineligible operators) on the fly. As soon as one operator isn't _is_blocked_on_downstream, return the usage (we can update this incrementally).

I recommend this approach because thinking about ineligble and eligble ops makes it hard to read. We can then remove the first if statement in _is_blocked_on_downstream

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

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

Thanks for your advice, I revised.

if self._op_resource_allocator is not None
else None
)
return budget is None or op.incremental_resource_usage().satisfies_limit(

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 help me understand the 2nd part of the if statement. Why would we consider the operator blocked if it can submit another task within its budget?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

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

We checked in_task_submission_backpressure first. The second part checks CPU, GPU, and memory:

  • If they fit, what's holding it back is object store memory, i.e., outputs that haven't been read yet.
  • If they don't fit, it's waiting for CPU, GPU, or memory that new nodes would add, so it still counts.

I reconstructed the logic to make it clearer.

Signed-off-by: Li Yonghui <yonghuili@google.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 and found 1 potential issue.

Fix All in Cursor

Reviewed by Cursor Bugbot for commit c12bbd8. Configure here.

Comment thread python/ray/data/_internal/execution/resource_manager.py
Signed-off-by: Li Yonghui <yonghuili@google.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

community-contribution Contributed by the community data Ray Data-related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Data] Ray Data continues autoscaling even when pipeline is backpressured by iteration

2 participants