Conversation
Signed-off-by: Sirui Huang <ray.huang@anyscale.com>
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Critical correctness issues and additional queue-accounting defects remain unresolved.
Review effort: Lite
Findings: 2
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.
Signed-off-by: Sirui Huang <ray.huang@anyscale.com>
|
|
||
| self._hash_linked.get_next() | ||
| self._hash_linked.add(first_bundle) | ||
| num_bundles_skipped += 1 |
There was a problem hiding this comment.
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)
Reviewed by Cursor Bugbot for commit d030f74. Configure here.
There was a problem hiding this comment.
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.
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>
There was a problem hiding this comment.
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).
Reviewed by Cursor Bugbot for commit 4bc13b5. Configure here.
4bc13b5 to
6fd5ea0
Compare
…once Signed-off-by: Sirui Huang <ray.huang@anyscale.com>
Signed-off-by: Sirui Huang <ray.huang@anyscale.com>




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.DefaultRankerorders 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 newhas_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
HashLinkedQueuewhenpreserve_orderis set orRAY_DATA_ENABLE_RESIDENT_FIRST_BUNDLE_QUEUES=0. The ranker is also on by default and falls back toDefaultRankerwhenRAY_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) andtest_zip.py(25 passed). Resident-input ranker is tested withpytest python/ray/data/tests/unit/test_ranker.py(3 passed) andtest_streaming_executor.py. No open PR covers this. AI assistance was used, but fully proofread and manually reviewed afterwards.