Repository navigation
[train] Pin worker groups to the reservation snapshot that sized them - #66646
Conversation
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Signed-off-by: Justin Yu <justinvyu@anyscale.com>
There was a problem hiding this comment.
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.
| 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 |
There was a problem hiding this comment.
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.
| 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}" | |
| ) |
Signed-off-by: Mark Towers <mark@anyscale.com>
Description
test_elastic_e2ehas 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 separaterecompute=Truepoll in_try_get_selectors, so:max_workersbundles, so afterResizeDecision(14)the reservation can reach 32 slots.len(label_selectors) != num_workersthen never passes.This PR takes the decision and the pins from the same snapshot, so they can't disagree:
ResizeDecision.ResizeDecision.label_selectorsis built by the scaling policy from the reservation that chosenum_workers, using the same_floor_slotsmath. 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.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.make_decision_for_non_running_worker_group(returningNoopDecisionuntil every worker is covered), the same way elastic waits formin_workers. This removesget_reserved_bundle_label_selectors, its 60s poll loop, and the controller'sWorkerGroupStartupTimeoutErrorraise. Both policies log a rate-limited "Waiting for reserved resources" message while they wait. Before, elastic logged nothing when nothing was reserved.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, soWorkerGroupnow pollspg_handle.wait()in 1s intervals and stops as soon as a pinned node is no longer alive. It raises the same retryableWorkerGroupStartupTimeoutErrorinstead of stalling for 60s.test_elastic_e2ehits 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 intoScalingPolicy._should_pin_to_reservation. The controller still lets a callback-providedlabel_selectorwin, 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_e2eflake differently: it picks a subset of a larger reservation with_compute_reservationsand adds a relative_FIT_EPSILONto the coordinator's shared_bundle_can_fit_on_node. This PR avoids the Train → Ray Data_compute_reservationsdependency 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):In the last
test_elastic_e2erun, 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
WorkerGroupStartupTimeoutErrorstill 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