Skip to content

[Data] Residency-aware Bundle Queue and Ranker Implementation - #66534

Open
rayhhome wants to merge 10 commits into
ray-project:masterfrom
rayhhome:bundle-queue
Open

rayhhome wants to merge 10 commits into
ray-project:masterfrom
rayhhome:bundle-queue

Conversation

@rayhhome

@rayhhome rayhhome commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

Description

This PR adds ResidentFirstBundleQueue, an input queue prioritizing bundles whose blocks are still in the object store over bundles whose blocks were evicted or stranded on a draining node.

The current linked-list queue may hand a lost bundle out, stalling the task on lineage reconstruction; the resident-first queue in this PR rotates such bundles to the back of the queue to overlap the reconstruction with useful work. It also recalculates its byte estimate from actual object locations on a fixed interval (default to 30 seconds), tracking queued memory that drives backpressure based on what the object store physically holds.

We also add ResidentInputRanker, which the streaming executor may use to pick the next operator to run. DefaultRanker orders eligible operators by whether they can be throttled and then by object store usage, so an operator whose next input bundle is lost can be chosen over one that could start a task immediately, and the task stalls on reconstruction. The new ranker additionally checks, through a new has_resident_next() query on the queue, whether any of the operator's input queues has a bundle whose blocks are in the object store. Operators that can run now are preferred, and ties still fall through to object store usage.

The queue is on by default and falls back to HashLinkedQueue when preserve_order is set or RAY_DATA_ENABLE_RESIDENT_FIRST_BUNDLE_QUEUES=0. The ranker is also on by default and falls back to DefaultRanker when RAY_DATA_USE_RESIDENT_INPUT_RANKER=0.

Related Issues/PRs

Prior attempt: #66537

Additional information

We've tried integrating the two components separately, but hit the wall when noticing that the ranker implementation is dependent on our resident-first bundle queue. Initial attempt here: #66537.

Resident-first bundle queue is tested with pytest python/ray/data/tests/test_bundle_queue.py (20 passed), test_streaming_executor.py (54 passed), test_actor_pool_map_operator.py (56 passed) and test_zip.py (25 passed). Resident-input ranker is tested with pytest python/ray/data/tests/unit/test_ranker.py (3 passed) and test_streaming_executor.py. No open PR covers this. AI assistance was used, but fully proofread and manually reviewed afterwards.

Signed-off-by: Sirui Huang <ray.huang@anyscale.com>
@rayhhome
rayhhome requested a review from a team as a code owner September 28, 2026 19:03
Copilot AI lite review requested due to automatic review settings September 28, 2026 19:03
@rayhhome rayhhome changed the title [Data] Node Aware ueue implementation Sep 28, 2026
@rayhhome rayhhome self-assigned this Sep 28, 2026
@rayhhome rayhhome added data Ray Data-related issues go add ONLY when ready to merge, run all tests labels Sep 28, 2026

@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 introduces ObjectStoreAwareBundleQueue, a thread-safe FIFO queue that prioritizes serving bundles whose blocks are still resident in the object store over those requiring lineage reconstruction. It also adds helper utilities to check object existence on non-drained nodes and comprehensive unit tests. The review feedback highlights several improvement opportunities: explicitly verifying that all block references are present in the object locations lookup to avoid false positives, returning the removed bundle in the remove method to adhere to the base class signature and Liskov Substitution Principle, adding the missing @override decorator to num_bundles, and using defensive .get() lookups to prevent potential KeyErrors when accessing object metadata.

Comment thread python/ray/data/_internal/utils/object_utils.py Outdated
Comment thread python/ray/data/_internal/execution/bundle_queue/object_store_aware.py Outdated
Comment thread python/ray/data/_internal/execution/bundle_queue/resident_first.py
Comment thread python/ray/data/_internal/utils/object_utils.py

@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/utils/cached_ray_internals.py
Comment thread python/ray/data/_internal/execution/bundle_queue/resident_first.py

Copilot AI 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.

Copilot review overview

🟡 Changes recommended

Critical correctness issues and additional queue-accounting defects remain unresolved.

Review effort: Lite
Findings: 2 High severity · 2 Medium severity

Open (4)
What changed in this PR

Adds a node-aware bundle queue that prioritizes resident object-store blocks and accounts for replicas and draining nodes.

Changes:

  • Adds object-store-aware queueing and selection logic.
  • Adds object-location and draining-node utilities.
  • Adds queue behavior, sizing, duplicate, and concurrency tests.
File Summary
python/​ray/​data/​tests/​test_bundle_queue.py Tests queue selection and behavior.
python/​ray/​data/​_internal/​utils/​object_utils.py Checks object residency.
python/​ray/​data/​_internal/​utils/​cached_ray_internals.py Identifies draining nodes.
python/​ray/​data/​_internal/​execution/​bundle_queue/​object_store_aware.py Implements node-aware queueing and size tracking.
python/​ray/​data/​_internal/​execution/​bundle_queue/​__init__.py Selects the configured bundle queue.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread python/ray/data/_internal/execution/bundle_queue/__init__.py Outdated
Comment thread python/ray/data/_internal/utils/object_utils.py
Comment thread python/ray/data/_internal/execution/bundle_queue/object_store_aware.py Outdated
Comment thread python/ray/data/_internal/utils/cached_ray_internals.py Outdated
@rayhhome rayhhome changed the title [Data] Node Aware Bundle Queue Implementation Sep 28, 2026
Signed-off-by: Sirui Huang <ray.huang@anyscale.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.


self._hash_linked.get_next()
self._hash_linked.add(first_bundle)
num_bundles_skipped += 1

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Lost bundles starve behind new arrivals

Medium Severity

_try_ensure_first_bundle_exists always rotates a missing bundle behind whatever is already queued, including bundles that arrive later. Combined with _refresh_bundle_sizes treating those objects as zero bytes, object-store backpressure no longer limits enqueue, so under eviction or drain a lost bundle can sit at the tail for the rest of a long or streaming job and never start reconstruction.

Additional Locations (1)
Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit d030f74. Configure here.

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.

Lost bundles are indeed starved behind new arrivals, but we have lineage reconstruction to bring blocks back. Oftentimes by the time when the queue sees the block is missing, the rebuild is already in flight. The only case that this would matter is where we have a bundle completely lost, i.e. reconstruction failed because lineage was evicted, retires were exhausted, or max_retires=0, but this is a very rare case. Great suggestion but I think we're trading runtime benefit for a fail-fast that guards a case that rarely happens.

Comment thread python/ray/data/_internal/execution/bundle_queue/resident_first.py Outdated
Signed-off-by: Sirui Huang <ray.huang@anyscale.com>
Signed-off-by: Sirui Huang <ray.huang@anyscale.com>
Signed-off-by: Sirui Huang <ray.huang@anyscale.com>
@rayhhome rayhhome changed the title [Data] Object-store-aware Bundle Queue Implementation Sep 29, 2026

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

There are 3 total unresolved issues (including 2 from previous reviews).

Fix All in Cursor

Reviewed by Cursor Bugbot for commit 4bc13b5. Configure here.

Comment thread python/ray/data/_internal/execution/bundle_queue/__init__.py Outdated
…once

Signed-off-by: Sirui Huang <ray.huang@anyscale.com>
Signed-off-by: Sirui Huang <ray.huang@anyscale.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

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

3 participants