Repository navigation
[Data] Suppress cluster scale-up when pipeline is backpressured by iteration - #66569
Yonghui-Lee wants to merge 6 commits into
Conversation
There was a problem hiding this comment.
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>
bc35ce8 to
5e5200a
Compare
|
Hi @Yonghui-Lee, thanks for the PR! I read the description and took a look at the PR and have 2 questions:
|
…caler-iteration-backpressure Signed-off-by: Li Yonghui <yonghuili@google.com>
|
@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):
|
Signed-off-by: Li Yonghui <yonghuili@google.com>
Signed-off-by: Li Yonghui <yonghuili@google.com>
iamjustinhsu
left a comment
There was a problem hiding this comment.
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() |
There was a problem hiding this comment.
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
There was a problem hiding this comment.
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( |
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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>
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes and found 1 potential issue.
Reviewed by Cursor Bugbot for commit c12bbd8. Configure here.
Signed-off-by: Li Yonghui <yonghuili@google.com>

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 bothRateBasedClusterAutoscalerandDefaultClusterAutoscalerV2.Root cause
Both autoscalers decide whether to scale up from
RollingLogicalUtilizationGauge, a rolling average ofglobal_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: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:
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 downstreamLimitleft out), and GPU-bound (downstream GPU op's usage kept).test_autoscaler_iteration_backpressure.py(new, registered inBUILD.bazel):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.