Skip to content

[train] Pin worker groups to the reservation snapshot that sized them - #66646

Merged
matthewdeng merged 3 commits into
ray-project:masterfrom
justinvyu:justinvyu/train-pins-on-decision
Oct 9, 2026
Merged

matthewdeng merged 3 commits into
ray-project:masterfrom
justinvyu:justinvyu/train-pins-on-decision

Conversation

@justinvyu

Copy link
Copy Markdown
Contributor

Description

test_elastic_e2e has been timing out on CI since #66002, which pins Train V2 workers to AutoscalingCoordinator reservations. The elastic worker group start waits the full 60s (RAY_TRAIN_WORKER_GROUP_START_TIMEOUT_S) whenever the reservation changes between the scaling decision and pinning. The decision is sized from the cached reservation, while the pins come from a separate recompute=True poll in _try_get_selectors, so:

  1. The reservation grows past the decision. Elastic requests max_workers bundles, so after ResizeDecision(14) the reservation can reach 32 slots. len(label_selectors) != num_workers then never passes.
  2. The reservation shrinks below the decision. After a node dies, the cached view still counts it, and pinning waits for workers that don't exist.

This PR takes the decision and the pins from the same snapshot, so they can't disagree:

  • Pins on ResizeDecision. ResizeDecision.label_selectors is built by the scaling policy from the reservation that chose num_workers, using the same _floor_slots math. The STRICT_PACK/STRICT_SPREAD layout check from [train] Pin Train V2 workers to autoscaling coordinator reservations … #66002 is kept. The controller just applies the pins and no longer polls the coordinator.
  • Elastic decides and pins from a fresh (recompute=True) view. In the non-running path it reads the fresh view only once the cached one looks ready, so waiting doesn't trigger a full coordinator tick every loop.
  • Fixed now waits for its reservation in make_decision_for_non_running_worker_group (returning NoopDecision until every worker is covered), the same way elastic waits for min_workers. This removes get_reserved_bundle_label_selectors, its 60s poll loop, and the controller's WorkerGroupStartupTimeoutError raise. Both policies log a rate-limited "Waiting for reserved resources" message while they wait. Before, elastic logged nothing when nothing was reserved.
  • Stop the placement group wait early when a pinned node dies. The coordinator drops nodes that ray.nodes() reports dead, but the GCS can still report a just-removed node as alive for a short time. A pin made in that window can never be placed, so WorkerGroup now polls pg_handle.wait() in 1s intervals and stops as soon as a pinned node is no longer alive. It raises the same retryable WorkerGroupStartupTimeoutError instead of stalling for 60s. test_elastic_e2e hits this window on purpose by removing every node and restarting right away. Locally this went from a 60s stall to giving up after about 1s.

Skipping pins (TPU, zero-resource workers, user label_selector) moves into ScalingPolicy._should_pin_to_reservation. The controller still lets a callback-provided label_selector win, because the policy can't see those.

Behavior change: a fixed run waiting for capacity no longer fails and retries every 60s through the failure policy. It waits in the controller loop like elastic. With the default controller_failure_limit=-1, the only visible difference is the log messages.

Related issues

Alternative to #66601, which fixes the same test_elastic_e2e flake differently: it picks a subset of a larger reservation with _compute_reservations and adds a relative _FIT_EPSILON to the coordinator's shared _bundle_can_fit_on_node. This PR avoids the Train → Ray Data _compute_reservations dependency and leaves the coordinator's fit check, which every requester shares, unchanged. It also avoids an extra restart when the reservation grows, because it starts at the size the reservation supports now.

Testing

All passing locally (macOS, ray_dev_py311):

pytest python/ray/train/v2/tests/test_autoscaling_coordinator_client.py \
       python/ray/train/v2/tests/test_elastic_scaling_policy.py              # 37 passed
pytest python/ray/train/v2/tests/test_controller.py \
       python/ray/train/v2/tests/test_controller_callback_behaviour.py \
       python/ray/train/v2/tests/test_failure_policy.py \
       python/ray/train/tests/test_autoscaling_coordinator_client.py         # 65 passed
pytest python/ray/train/v2/tests/test_scheduling.py \
       python/ray/train/v2/tests/test_elastic_e2e.py \
       python/ray/train/v2/tests/test_state.py                               # 34 passed
pytest python/ray/train/v2/tests/test_worker_group.py \
       python/ray/train/v2/tests/test_elastic_e2e.py                         # 43 passed
pre-commit run --files <changed files>                                       # clean

In the last test_elastic_e2e run, the controller log shows the "decide 4 right after every node was removed" restart stopping after about 1s ("Giving up on the placement group early…") instead of the 60s stall seen before the early-stop change.

Known gap: when it stops early, the retried WorkerGroupStartupTimeoutError still says "timed out after 60.0 seconds". The info log just before it gives the real reason.

AI assistance (Claude Code) was used to write this change.

🤖 Generated with Claude Code

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Signed-off-by: Justin Yu <justinvyu@anyscale.com>

@pseudo-rnd-thoughts pseudo-rnd-thoughts left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Looks good

@pseudo-rnd-thoughts pseudo-rnd-thoughts added train Ray Train Related Issue core-autoscaler autoscaler related issues go add ONLY when ready to merge, run all tests labels Oct 2, 2026
@pseudo-rnd-thoughts
pseudo-rnd-thoughts marked this pull request as ready for review October 2, 2026 16:02
@pseudo-rnd-thoughts
pseudo-rnd-thoughts requested a review from a team as a code owner October 2, 2026 16:02

@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 refactors the scaling policy and worker group initialization in Ray Train v2 to handle node-pinning and reservation decisions more robustly. It updates both elastic and fixed scaling policies to include node-pinning label selectors directly in the resize decisions, built from fresh reservation snapshots. Additionally, the worker group now monitors pinned nodes and aborts waiting for placement groups early if a pinned node dies. The review feedback suggests improving robustness against transient failures: wrapping the ray.nodes() liveness check in a try-except block to handle GCS query failures, and explicitly handling None returns from _get_reserved_resources(recompute=True) in both elastic and fixed scaling policies to prevent misleading log outputs.

Comment on lines +290 to +297
alive_node_ids = {node["NodeID"] for node in ray.nodes() if node["Alive"]}
dead_node_ids = pinned_node_ids - alive_node_ids
if dead_node_ids:
logger.info(
"Giving up on the placement group early: the worker group is "
f"pinned to nodes that are no longer alive: {sorted(dead_node_ids)}."
)
return False

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.

medium

Querying ray.nodes() directly inside the placement group wait loop can raise transient exceptions (e.g., RaySystemError or gRPC timeouts) if the GCS is temporarily overloaded or recovering from a node failure. Wrapping this call in a try-except block ensures that transient GCS query failures do not crash the entire controller/worker group startup, allowing it to retry or continue waiting gracefully.

Suggested change
alive_node_ids = {node["NodeID"] for node in ray.nodes() if node["Alive"]}
dead_node_ids = pinned_node_ids - alive_node_ids
if dead_node_ids:
logger.info(
"Giving up on the placement group early: the worker group is "
f"pinned to nodes that are no longer alive: {sorted(dead_node_ids)}."
)
return False
try:
alive_node_ids = {node["NodeID"] for node in ray.nodes() if node["Alive"]}
dead_node_ids = pinned_node_ids - alive_node_ids
if dead_node_ids:
logger.info(
"Giving up on the placement group early: the worker group is "
f"pinned to nodes that are no longer alive: {sorted(dead_node_ids)}."
)
return False
except Exception as e:
logger.warning(
f"Failed to query ray.nodes() to check pinned node liveness: {e}"
)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Fixed this issue in 53e5345

Comment thread python/ray/train/v2/_internal/execution/scaling_policy/elastic.py
Comment thread python/ray/train/v2/_internal/execution/scaling_policy/fixed.py
Mark Towers and others added 2 commits October 9, 2026 10:51
@matthewdeng
matthewdeng merged commit c854475 into ray-project:master Oct 9, 2026
6 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

core-autoscaler autoscaler related issues go add ONLY when ready to merge, run all tests train Ray Train Related Issue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants