[Data] Fix misleading Shuffle v2 progress bars for logs - #66456
machichima wants to merge 10 commits into
Conversation
Signed-off-by: machichima <nary12321@gmail.com>
Signed-off-by: machichima <nary12321@gmail.com>
Signed-off-by: machichima <nary12321@gmail.com>
Signed-off-by: machichima <nary12321@gmail.com>
Signed-off-by: machichima <nary12321@gmail.com>
Signed-off-by: machichima <nary12321@gmail.com>
…-v2-progress-bar Signed-off-by: machichima <nary12321@gmail.com> # Conflicts: # python/ray/data/_internal/execution/operators/shuffle_operators/disk_shuffle_map_operator.py # python/ray/data/_internal/execution/operators/shuffle_operators/disk_shuffle_reduce_operator.py # python/ray/data/tests/test_hash_shuffle_v2.py
|
@owenowenisme PTAL when you have time. Thank you! |
There was a problem hiding this comment.
Code Review
This pull request removes sub-progress bars from the shuffle operators and refactors how the total number of output rows is estimated. Specifically, it introduces a preserves_row_count flag for reduce operators and ensures that map operators with block transformers return None for their total row count, as the count is unknown until execution. In the feedback, a reviewer pointed out that removing the 'or 1' fallback in logging_progress.py might cause empty datasets (with a total of 0) to display as '0/?' instead of '0/0' due to a falsy check in _format_progress, and suggested explicitly checking for None instead.
Signed-off-by: machichima <nary12321@gmail.com>
Signed-off-by: machichima <nary12321@gmail.com>
| for op_state, (tid, progress, _) in self._op_display.items(): | ||
| completed = op_state.op.metrics.row_outputs_taken |
There was a problem hiding this comment.
After the change, I saw 0/0 in rich view when completed:
❯ RAY_DATA_ENABLE_RICH_PROGRESS_BARS=1 RAY_TQDM=0 python shuffle-progress-bar.py
2026-09-27 16:06:53,261 WARNING authentication_token_setup.py:85 -- Token authentication is enabled for this Ray cluster. Set RAY_AUTH_MODE=disabled to opt out. For more information, see https://docs.ray.io/en/latest/ray-security/token-auth.html
2026-09-27 16:06:56,887 INFO worker.py:2019 -- Started a local Ray instance. View the dashboard at http://127.0.0.1:8265
2026-09-27 16:06:58,023 INFO streaming_executor.py:215 -- Starting execution of Dataset dataset_034VykfeDKyOn7QgJUM38f_0. Full logs are in /tmp/ray/session_2026-09-27_16-06-53_262122_2903/logs/ray-data
2026-09-27 16:06:58,023 INFO streaming_executor.py:216 -- Execution plan of Dataset dataset_034VykfeDKyOn7QgJUM38f_0: InputDataBuffer[Input] -> TaskPoolMapOperator[ReadRange->MapBatches(add_column)] -> ShuffleMapOp[HashAggregateMap(key_columns=('keys',), num_partitions=200)] -> ShuffleReduceOp[HashAggregateReduce(key_columns=('keys',), num_partitions=200)]
• ✔️ Dataset dataset_034VykfeDKyOn7QgJUM38f_0 execution finished in 9.57 seconds 100% ━━━━━━━━━━━━━━━ 10/? [ 0:00:09 , ? rows/s ]
• ✔️ Dataset dataset_034VykfeDKyOn7QgJUM38f_0 execution finished in 9.57 seconds 100% ━━━━━━━━━━━━━━━ 10/? [ 0:00:09 , ? rows/s ]
│ Active/total resources: Active & requested resources: 0/16 CPU, 160.0B/1.0GiB object store
│
├─ ReadRange->MapBatches(add_column) 100% ━━━━━━━━━━ 100000k/100000k [ 0:00:07 , 13468.34 k row/s ]
│ Tasks: 0; Actors: 0; Queued blocks: 0 (0.0B); Resources: 0.0 CPU, 0.0B object store
├─ ⠇ HashAggregateMap(key_columns=('keys',), num_partitions=200) 0% ━━━━━━━━━━ 0/0 [ 0:00:09 , ? rows/s ]
│ Tasks: 0; Actors: 0; Queued blocks: 0 (0.0B); Resources: 0.0 CPU, 0.0B object store; map: 2/2
├─ ⠇ HashAggregateReduce(key_columns=('keys',), num_partitions=200) 0% ━━━━━━━━━━ 0/0 [ 0:00:09 , ? rows/s ]
│ Tasks: 0; Actors: 0; Queued blocks: 0 (0.0B); Resources: 0.0 CPU, 160.0B object store; reduce: 10/10
2026-09-27 16:07:07,621 INFO streaming_executor.py:371 -- ✔️ Dataset dataset_034VykfeDKyOn7QgJUM38f_0 execution finished in 9.57 seconds
It's because completed will be set to 0 when num_output_rows_total returns None (updated in this PR):
ray/python/ray/data/_internal/progress/rich_progress.py
Lines 387 to 390 in f239110
After this update, the rich view can show corretly:
# In progress
⠧ Dataset dataset_034VznJ8Pa8TX9hTpMVvuw_0 running: 0% ━━━━━━━━━━━━━━━ 0/? [ 0:00:08 , ? rows/s ]
│ Active/total resources: Active & requested resources: 10/16 CPU, 640.0B/32.8GiB memory, 0.0B/1.0GiB object store
│
├─ ReadRange->MapBatches(add_column) 100% ━━━━━━━━━━ 100000k/100000k [ 0:00:07 , 13104.78 k row/s ]
│ Tasks: 0; Actors: 0; Queued blocks: 0 (0.0B); Resources: 0.0 CPU, 0.0B object store
├─ ⠧ HashAggregateMap(key_columns=('keys',), num_partitions=200) 0% ━━━━━━━━━━ 20/? [ 0:00:08 , ? rows/s ]
│ Tasks: 0; Actors: 0; Queued blocks: 0 (0.0B); Resources: 0.0 CPU, 0.0B object store; map: 2/2
├─ ⠧ HashAggregateReduce(key_columns=('keys',), num_partitions=200) 0% ━━━━━━━━━━ 0/? [ 0:00:08 , ? rows/s ]
│ Tasks: 10; Actors: 0; Queued blocks: 0 (0.0B); Resources: 10.0 CPU, 640.0B memory, 0.0B object store; reduce: 0/10
# Completed
• ✔️ Dataset dataset_034Vze2rD7OxlgNqnBiU3M_0 execution finished in 10.09 seconds 100% ━━━━━━━━━━━━━━━ 10/? [ 0:00:09 , ? rows/s ]
• ✔️ Dataset dataset_034Vze2rD7OxlgNqnBiU3M_0 execution finished in 10.09 seconds 100% ━━━━━━━━━━━━━━━ 10/? [ 0:00:09 , ? rows/s ]
│ Active/total resources: Active & requested resources: 0/16 CPU, 160.0B/1.0GiB object store
│
├─ ReadRange->MapBatches(add_column) 100% ━━━━━━━━━━ 100000k/100000k [ 0:00:07 , 12542.56 k row/s ]
│ Tasks: 0; Actors: 0; Queued blocks: 0 (0.0B); Resources: 0.0 CPU, 0.0B object store
├─ HashAggregateMap(key_columns=('keys',), num_partitions=200) 100% ━━━━━━━━━━ 20/20 [ 0:00:09 , 2.01 row/s ]
│ Tasks: 0; Actors: 0; Queued blocks: 0 (0.0B); Resources: 0.0 CPU, 0.0B object store; map: 2/2
├─ HashAggregateReduce(key_columns=('keys',), num_partitions=200) 100% ━━━━━━━━━━ 10/10 [ 0:00:09 , 1.00 row/s ]
│ Tasks: 0; Actors: 0; Queued blocks: 0 (0.0B); Resources: 0.0 CPU, 160.0B object store; reduce: 10/10
2026-09-27 16:43:31,746 INFO streaming_executor.py:371 -- ✔️ Dataset dataset_034Vze2rD7OxlgNqnBiU3M_0 execution finished in 10.09 seconds
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes using default effort and found 1 potential issue.
Reviewed by Cursor Bugbot for commit 7259ee4. Configure here.
|
|
||
|
|
||
| class ShuffleMapOp(InternalQueueOperatorMixin, PhysicalOperator, SubProgressBarMixin): | ||
| class ShuffleMapOp(InternalQueueOperatorMixin, PhysicalOperator): |
There was a problem hiding this comment.
Shuffle v2 bars hidden without verbose
Medium Severity
Dropping SubProgressBarMixin from the v2 shuffle ops also drops their operator-level bars when verbose_progress is off. Progress managers treat that mixin as the AllToAll marker, so RAY_DATA_VERBOSE_PROGRESS=0 now hides HashShuffle* / HashAggregate* entirely instead of only the nested Map/Reduce lines.
Additional Locations (2)
Reviewed by Cursor Bugbot for commit 7259ee4. Configure here.


Description
Current Shuffle v2 (ShuffleStrategy.SHUFFLE_V2) progress bars are misleading:
HashAggregateReduce(key_columns=('keys',), num_partitions=200): 3/6000000 - Reduce: 3/6000000Main changes
preserves_row_countargs to reduce op to identify aggregated/non-aggregated operations. Useblock_transformerto identify for map op.num_output_rows_total()returnNoneifpreserves_row_countis false orblock_transformerexistsself.input_dependencies[0].num_output_rows_total()to get the output rows from input directlyor 1default inlogging_progress.pywhen initializing an operator's total. Without this, the log rendered an unknown total as10/1.Related issues
Closes #66314
Additional information
Perform manual tests for aggregated and non-aggregated operations, and check the log output:
Aggregated operation
RAY_DATA_NON_TTY_PROGRESS_LOG_INTERVAL=0 python shuffle-progress-bar.py > shuffle-logs.txt 2>&1?for the total instead of the wrong valueNon-aggregated operation
RAY_DATA_NON_TTY_PROGRESS_LOG_INTERVAL=0 python repartition-progress-bar.py > logs-repartition.txt 2>&1100M) when operation not finished