From d310dbe401b9f8a23c5b6f6cd787ef7935fc9078 Mon Sep 17 00:00:00 2001 From: machichima Date: Thu, 24 Sep 2026 21:01:46 +0800 Subject: [PATCH 1/8] feat: remove subprogress bar Signed-off-by: machichima --- .../external_shuffle_map_operator.py | 30 +------------------ .../external_shuffle_reduce_operator.py | 21 +------------ .../shuffle_operators/shuffle_map_operator.py | 30 +------------------ .../shuffle_reduce_operator.py | 21 +------------ 4 files changed, 4 insertions(+), 98 deletions(-) diff --git a/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_map_operator.py b/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_map_operator.py index 6a1304855d0b..7e84c91394f4 100644 --- a/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_map_operator.py +++ b/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_map_operator.py @@ -21,7 +21,6 @@ from ray.data._internal.execution.interfaces.physical_operator import ( MetadataOpTask, OpTask, - estimate_total_num_of_blocks, ) from ray.data._internal.execution.operators.base_physical_operator import ( InternalQueueOperatorMixin, @@ -37,7 +36,6 @@ from ray.data._internal.execution.operators.shuffle_operators.shuffle_tasks import ( SHUFFLE_PEAK_MEMORY_MULTIPLIER, ) -from ray.data._internal.execution.operators.sub_progress import SubProgressBarMixin from ray.data.block import BlockExecStats, BlockMetadata, BlockStats from ray.data.context import DataContext from ray.types import ObjectRef @@ -46,8 +44,6 @@ if typing.TYPE_CHECKING: import pyarrow as pa - from ray.data._internal.progress.base_progress import BaseProgressBar - logger = logging.getLogger(__name__) @@ -58,9 +54,7 @@ def _make_mapper_sentinel(mapper_id: int) -> Tuple[str, ...]: return (f"{_MAPPER_ID_SENTINEL}{mapper_id}",) -class ExternalHashShuffleMapOp( - InternalQueueOperatorMixin, PhysicalOperator, SubProgressBarMixin -): +class ExternalHashShuffleMapOp(InternalQueueOperatorMixin, PhysicalOperator): """External-shuffle map operator. See module docstring.""" _DEFAULT_SHUFFLE_MAP_TASK_NUM_CPUS = 1.0 @@ -124,9 +118,6 @@ def __init__( self._partition_rows: Dict[int, int] = defaultdict(int) self._partition_bytes: Dict[int, int] = defaultdict(int) - # -- Sub-progress bars ----------------------------------------------- - self._map_bar: Optional["BaseProgressBar"] = None - # ===================================================================== # External-shuffle-specific state below. # ===================================================================== @@ -274,15 +265,6 @@ def _submit_shuffle_map_task( task_id=task.get_task_id(), ) - if self._map_bar is not None: - _, _, num_rows = estimate_total_num_of_blocks( - cur_task_idx + 1, - self.upstream_op_num_outputs(), - self._metrics, - total_num_tasks=None, - ) - self._map_bar.update(total=num_rows) - def _handle_map_done( self, task_idx: int, @@ -365,9 +347,6 @@ def _handle_map_done( task_exec_driver_stats=None, ) - if self._map_bar is not None: - self._map_bar.update(increment=input_rows) - self._maybe_emit_partition_bundles() def _maybe_emit_partition_bundles(self) -> None: @@ -567,13 +546,6 @@ def progress_str(self) -> str: parts.append(f"merge_buf: {total_merge_buf}") return ", ".join(parts) - def get_sub_progress_bar_names(self) -> Optional[List[str]]: - return ["Map"] - - def set_sub_progress_bar(self, name: str, pg: "BaseProgressBar") -> None: - if name == "Map": - self._map_bar = pg - @property def num_partitions(self) -> int: return self._num_partitions diff --git a/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_reduce_operator.py b/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_reduce_operator.py index a5e2faeb4738..f2e988714194 100644 --- a/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_reduce_operator.py +++ b/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_reduce_operator.py @@ -39,7 +39,6 @@ from ray.data._internal.execution.operators.shuffle_operators.shuffle_tasks import ( SHUFFLE_PEAK_MEMORY_MULTIPLIER, ) -from ray.data._internal.execution.operators.sub_progress import SubProgressBarMixin from ray.data.block import BlockAccessor, BlockStats, TaskExecWorkerStats, to_stats from ray.data.context import DataContext @@ -47,12 +46,11 @@ from ray.data._internal.execution.operators.map_transformer import ( MapTransformer, ) - from ray.data._internal.progress.base_progress import BaseProgressBar logger = logging.getLogger(__name__) -class ExternalHashShuffleReduceOp(PhysicalOperator, SubProgressBarMixin): +class ExternalHashShuffleReduceOp(PhysicalOperator): """External-shuffle reduce operator. Structurally mirrors ``ShuffleReduceOp``: one wrapper bundle per partition @@ -132,9 +130,6 @@ def __init__( # -- Stats ----------------------------------------------------------- self._output_blocks_stats: List[BlockStats] = [] - # -- Sub-progress bars ----------------------------------------------- - self._reduce_bar: Optional["BaseProgressBar"] = None - # ===================================================================== # External-shuffle-specific state below. # ===================================================================== @@ -318,8 +313,6 @@ def _emit_empty_partition(self, refs: RefBundle, schema: pa.Schema) -> None: ) self._estimated_num_output_bundles = num_outputs self._estimated_output_num_rows = num_rows - if self._reduce_bar is not None: - self._reduce_bar.update(increment=0, total=self.num_output_rows_total()) def has_next(self) -> bool: return len(self._output_queue) > 0 @@ -345,11 +338,6 @@ def _handle_reduce_output_ready(self, partition_id: int, bundle: RefBundle) -> N ) self._estimated_num_output_bundles = num_outputs self._estimated_output_num_rows = num_rows - if self._reduce_bar is not None: - self._reduce_bar.update( - increment=bundle.num_rows() or 0, - total=self.num_output_rows_total(), - ) def _handle_reduce_done( self, @@ -441,10 +429,3 @@ def progress_str(self) -> str: submitted = self._num_reduce_tasks_submitted done = submitted - len(self._shuffle_reduce_tasks) return f"reduce: {done}/{submitted}" - - def get_sub_progress_bar_names(self) -> Optional[List[str]]: - return ["Reduce"] - - def set_sub_progress_bar(self, name: str, pg: "BaseProgressBar") -> None: - if name == "Reduce": - self._reduce_bar = pg diff --git a/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_map_operator.py b/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_map_operator.py index 2b460b738209..1bef71b2628e 100644 --- a/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_map_operator.py +++ b/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_map_operator.py @@ -1,7 +1,6 @@ import dataclasses import functools import logging -import typing from collections import defaultdict from typing import Any, Dict, List, Optional, Tuple @@ -19,7 +18,6 @@ from ray.data._internal.execution.interfaces.physical_operator import ( MetadataOpTask, OpTask, - estimate_total_num_of_blocks, ) from ray.data._internal.execution.operators.base_physical_operator import ( InternalQueueOperatorMixin, @@ -30,15 +28,11 @@ PartitionFn, _shuffle_map_task, ) -from ray.data._internal.execution.operators.sub_progress import SubProgressBarMixin from ray.data.block import Block, BlockMetadata, BlockStats from ray.data.context import DataContext from ray.types import ObjectRef from ray.util.scheduling_strategies import NodeAffinitySchedulingStrategy -if typing.TYPE_CHECKING: - from ray.data._internal.progress.base_progress import BaseProgressBar - logger = logging.getLogger(__name__) @@ -61,7 +55,7 @@ def extract_partition_id(bundle: RefBundle) -> int: raise ValueError("ShuffleMapOp bundle is missing a partition_id sentinel.") -class ShuffleMapOp(InternalQueueOperatorMixin, PhysicalOperator, SubProgressBarMixin): +class ShuffleMapOp(InternalQueueOperatorMixin, PhysicalOperator): """Map phase of a shuffle: partition inputs and group shards by partition. Each map task splits its input into num_partitions shards. Shards land in a @@ -149,9 +143,6 @@ def __init__( self._total_input_bytes: int = 0 self._map_blocks_stats: List[BlockStats] = [] - # -- Sub-progress bars ----------------------------------------------- - self._map_bar: Optional["BaseProgressBar"] = None - @property def _input_queues(self) -> List[BaseBundleQueue]: return [] @@ -268,15 +259,6 @@ def _submit_shuffle_map_task( task_id=task.get_task_id(), ) - if self._map_bar is not None: - _, _, num_rows = estimate_total_num_of_blocks( - cur_task_idx + 1, - self.upstream_op_num_outputs(), - self._metrics, - total_num_tasks=None, - ) - self._map_bar.update(total=num_rows) - def _handle_map_done( self, task_idx: int, @@ -322,9 +304,6 @@ def _handle_map_done( task_exec_driver_stats=None, ) - if self._map_bar is not None: - self._map_bar.update(increment=input_meta.num_rows or 0) - self._maybe_emit_partition_bundles() def _maybe_emit_partition_bundles(self) -> None: @@ -447,10 +426,3 @@ def progress_str(self) -> str: if total_merge_buf: parts.append(f"merge_buf: {total_merge_buf}") return ", ".join(parts) - - def get_sub_progress_bar_names(self) -> Optional[List[str]]: - return ["Map"] - - def set_sub_progress_bar(self, name: str, pg: "BaseProgressBar") -> None: - if name == "Map": - self._map_bar = pg diff --git a/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_reduce_operator.py b/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_reduce_operator.py index 670796fc99fc..540bdf6b6bc2 100644 --- a/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_reduce_operator.py +++ b/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_reduce_operator.py @@ -29,13 +29,11 @@ ReduceFn, _shuffle_reduce_task, ) -from ray.data._internal.execution.operators.sub_progress import SubProgressBarMixin from ray.data.block import BlockAccessor, BlockStats, TaskExecWorkerStats, to_stats from ray.data.context import DataContext if typing.TYPE_CHECKING: from ray.data._internal.execution.operators.map_transformer import MapTransformer - from ray.data._internal.progress.base_progress import BaseProgressBar logger = logging.getLogger(__name__) @@ -56,7 +54,7 @@ def _merged_reduce_runtime_env(user_runtime_env: Dict[str, Any]) -> Dict[str, An return merged -class ShuffleReduceOp(PhysicalOperator, SubProgressBarMixin): +class ShuffleReduceOp(PhysicalOperator): """Reduce phase of a shuffle. Supports one or more co-partitioned upstream `ShuffleMapOp`s. With a single @@ -158,9 +156,6 @@ def __init__( # -- Stats ----------------------------------------------------------- self._output_blocks_stats: List[BlockStats] = [] - # -- Sub-progress bars ----------------------------------------------- - self._reduce_bar: Optional["BaseProgressBar"] = None - def _reduce_task_remote_args(self, memory_estimate: int) -> Dict[str, Any]: remote_args: Dict[str, Any] = { "num_cpus": self._DEFAULT_SHUFFLE_REDUCE_TASK_NUM_CPUS, @@ -356,8 +351,6 @@ def _emit_empty_partition(self, refs: RefBundle, schema: pa.Schema) -> None: ) self._estimated_num_output_bundles = num_outputs self._estimated_output_num_rows = num_rows - if self._reduce_bar is not None: - self._reduce_bar.update(increment=0, total=self.num_output_rows_total()) def has_next(self) -> bool: return len(self._output_queue) > 0 @@ -383,11 +376,6 @@ def _handle_reduce_output_ready(self, partition_id: int, bundle: RefBundle) -> N ) self._estimated_num_output_bundles = num_outputs self._estimated_output_num_rows = num_rows - if self._reduce_bar is not None: - self._reduce_bar.update( - increment=bundle.num_rows() or 0, - total=self.num_output_rows_total(), - ) def _handle_reduce_done( self, @@ -479,10 +467,3 @@ def progress_str(self) -> str: submitted = self._num_reduce_tasks_submitted done = submitted - len(self._shuffle_reduce_tasks) return f"reduce: {done}/{submitted}" - - def get_sub_progress_bar_names(self) -> Optional[List[str]]: - return ["Reduce"] - - def set_sub_progress_bar(self, name: str, pg: "BaseProgressBar") -> None: - if name == "Reduce": - self._reduce_bar = pg From 8f1615af90ad59410d05432d432e000c78602672 Mon Sep 17 00:00:00 2001 From: machichima Date: Sat, 26 Sep 2026 21:08:32 +0800 Subject: [PATCH 2/8] feat: add preserves_row_count flag to reduce Signed-off-by: machichima --- .../shuffle_operators/external_shuffle_reduce_operator.py | 4 ++++ .../operators/shuffle_operators/shuffle_reduce_operator.py | 6 ++++++ python/ray/data/_internal/logical/rules/operator_fusion.py | 4 ++++ python/ray/data/_internal/planner/plan_all_to_all_op.py | 2 ++ 4 files changed, 16 insertions(+) diff --git a/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_reduce_operator.py b/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_reduce_operator.py index f2e988714194..e04ee41fb2e0 100644 --- a/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_reduce_operator.py +++ b/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_reduce_operator.py @@ -79,6 +79,7 @@ def __init__( peak_memory_multiplier: float = SHUFFLE_PEAK_MEMORY_MULTIPLIER, name: str = "ExternalHashShuffleReduce", should_emit_empty_partitions: bool = True, + preserves_row_count: bool = True, fused_output_map_transformer: Optional["MapTransformer"] = None, fused_output_map_task_kwargs: Optional[Dict[str, Any]] = None, fused_output_map_target_max_block_size_override: Optional[int] = None, @@ -103,6 +104,9 @@ def __init__( self._reduce_fn: ReduceFn = reduce_fn self._disallow_block_splitting: bool = disallow_block_splitting self._emit_empty_partitions: bool = should_emit_empty_partitions + # False when reduce_fn (aggregation) or a fused map can change the row + # count, so num_output_rows_total() can't borrow the map op's total. + self._preserves_row_count: bool = preserves_row_count self._peak_memory_multiplier: float = peak_memory_multiplier # -- Reduce task config & tracking ----------------------------------- diff --git a/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_reduce_operator.py b/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_reduce_operator.py index 540bdf6b6bc2..bc7c64e2a241 100644 --- a/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_reduce_operator.py +++ b/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_reduce_operator.py @@ -84,6 +84,8 @@ class ShuffleReduceOp(PhysicalOperator): name: Display name shown in progress bars and logs. should_emit_empty_partitions: If True (default), an empty partition emits one schema-only placeholder block. + preserves_row_count: If False, the reduce may change the row count + (aggregation, fused map), so the output row total is unknown. fused_output_map_transformer: Set by ``FuseOperators`` when a ``TaskPoolMapOperator`` directly downstream is fused into this reduce: each reduce task applies it to its output blocks before @@ -108,6 +110,7 @@ def __init__( peak_memory_multiplier: float = SHUFFLE_PEAK_MEMORY_MULTIPLIER, name: str = "ShuffleReduce", should_emit_empty_partitions: bool = True, + preserves_row_count: bool = True, fused_output_map_transformer: Optional["MapTransformer"] = None, fused_output_map_task_kwargs: Optional[Dict[str, Any]] = None, fused_output_map_target_max_block_size_override: Optional[int] = None, @@ -127,6 +130,9 @@ def __init__( self._reduce_fn: ReduceFn = reduce_fn self._disallow_block_splitting: bool = disallow_block_splitting self._emit_empty_partitions: bool = should_emit_empty_partitions + # False when reduce_fn (aggregation) or a fused map can change the row + # count, so num_output_rows_total() can't borrow the map op's total. + self._preserves_row_count: bool = preserves_row_count self._peak_memory_multiplier: float = peak_memory_multiplier # -- Reduce task config & tracking ----------------------------------- diff --git a/python/ray/data/_internal/logical/rules/operator_fusion.py b/python/ray/data/_internal/logical/rules/operator_fusion.py index 27fad2dd40d9..f51297161eb6 100644 --- a/python/ray/data/_internal/logical/rules/operator_fusion.py +++ b/python/ray/data/_internal/logical/rules/operator_fusion.py @@ -411,6 +411,8 @@ def _get_fused_map_into_shuffle_reduce_operator( reduce_ray_remote_args=up_op._reduce_ray_remote_args, peak_memory_multiplier=up_op._peak_memory_multiplier, should_emit_empty_partitions=up_op._emit_empty_partitions, + # The fused map (filter/flat_map/...) may change the row count. + preserves_row_count=False, name=name, fused_output_map_transformer=down_op.get_map_transformer(), fused_output_map_task_kwargs=down_op.get_map_task_kwargs(), @@ -428,6 +430,8 @@ def _get_fused_map_into_shuffle_reduce_operator( reduce_ray_remote_args=up_op._reduce_ray_remote_args, peak_memory_multiplier=up_op._peak_memory_multiplier, should_emit_empty_partitions=up_op._emit_empty_partitions, + # The fused map (filter/flat_map/...) may change the row count. + preserves_row_count=False, name=name, fused_output_map_transformer=down_op.get_map_transformer(), fused_output_map_task_kwargs=down_op.get_map_task_kwargs(), diff --git a/python/ray/data/_internal/planner/plan_all_to_all_op.py b/python/ray/data/_internal/planner/plan_all_to_all_op.py index cec45b629482..d7d1cd000afc 100644 --- a/python/ray/data/_internal/planner/plan_all_to_all_op.py +++ b/python/ray/data/_internal/planner/plan_all_to_all_op.py @@ -217,6 +217,8 @@ def _plan_hash_shuffle_aggregate_v2( # block; a placeholder would carry the map's pre-finalize schema and # conflict with finalized non-empty partitions. should_emit_empty_partitions=False, + # Aggregation collapses each key group to one row. + preserves_row_count=False, name=( f"{prefix}HashAggregateReduce(key_columns={key_columns}, " f"num_partitions={num_partitions})" From 8d39c9b16bc60df38fc49be1084ec0a5506f4b02 Mon Sep 17 00:00:00 2001 From: machichima Date: Sun, 27 Sep 2026 09:50:39 +0800 Subject: [PATCH 3/8] feat: num out rows return none if is aggregate Signed-off-by: machichima --- .../shuffle_operators/external_shuffle_map_operator.py | 5 +++++ .../shuffle_operators/external_shuffle_reduce_operator.py | 8 ++++---- .../operators/shuffle_operators/shuffle_map_operator.py | 5 +++++ .../shuffle_operators/shuffle_reduce_operator.py | 7 ++++--- 4 files changed, 18 insertions(+), 7 deletions(-) diff --git a/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_map_operator.py b/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_map_operator.py index 7e84c91394f4..b5daa2cee1bc 100644 --- a/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_map_operator.py +++ b/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_map_operator.py @@ -513,6 +513,11 @@ def get_stats(self) -> Dict[str, List[BlockStats]]: return {self._name: self._map_blocks_stats} def num_output_rows_total(self) -> Optional[int]: + # The aggregation combiner (block_transformer) pre-aggregates each map + # task's input before partitioning, so the output row count is unknown + # until the maps run. + if self._block_transformer is not None: + return None return self._total_input_rows if self._total_input_rows > 0 else None def current_logical_usage(self) -> ExecutionResources: diff --git a/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_reduce_operator.py b/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_reduce_operator.py index e04ee41fb2e0..ab91ecd245a1 100644 --- a/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_reduce_operator.py +++ b/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_reduce_operator.py @@ -394,10 +394,10 @@ def get_stats(self) -> Dict[str, List[BlockStats]]: return {self._name: self._output_blocks_stats} def num_output_rows_total(self) -> Optional[int]: - # Multi-input reduces (e.g. join) can grow or shrink the row count, so - # it is unknown until the reducers run; a single-input reduce preserves - # it. - if self._num_inputs > 1: + # Multi-input reduces (e.g. join) and non-row-preserving reduces + # (aggregation, fused map) can grow or shrink the row count, so it is + # unknown until the reducers run. + if self._num_inputs > 1 or not self._preserves_row_count: return None upstream = self.input_dependencies[0] assert isinstance(upstream, ExternalHashShuffleMapOp) diff --git a/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_map_operator.py b/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_map_operator.py index 1bef71b2628e..f44380bcba2a 100644 --- a/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_map_operator.py +++ b/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_map_operator.py @@ -393,6 +393,11 @@ def get_stats(self) -> Dict[str, List[BlockStats]]: return {self._name: self._map_blocks_stats} def num_output_rows_total(self) -> Optional[int]: + # The aggregation combiner (block_transformer) pre-aggregates each map + # task's input before partitioning, so the output row count is unknown + # until the maps run. + if self._block_transformer is not None: + return None return self._total_input_rows if self._total_input_rows > 0 else None def current_logical_usage(self) -> ExecutionResources: diff --git a/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_reduce_operator.py b/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_reduce_operator.py index bc7c64e2a241..a162b5f501c8 100644 --- a/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_reduce_operator.py +++ b/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_reduce_operator.py @@ -434,9 +434,10 @@ def get_stats(self) -> Dict[str, List[BlockStats]]: return {self._name: self._output_blocks_stats} def num_output_rows_total(self) -> Optional[int]: - # Multi-input reduces (e.g. join) can grow or shrink the row count, so it - # is unknown until the reducers run; a single-input reduce preserves it. - if self._num_inputs > 1: + # Multi-input reduces (e.g. join) and non-row-preserving reduces + # (aggregation, fused map) can grow or shrink the row count, so it is + # unknown until the reducers run. + if self._num_inputs > 1 or not self._preserves_row_count: return None upstream = self.input_dependencies[0] assert isinstance(upstream, ShuffleMapOp) From 060b86c6dadfef26943a71fa035c19942a2a6af3 Mon Sep 17 00:00:00 2001 From: machichima Date: Sun, 27 Sep 2026 09:51:35 +0800 Subject: [PATCH 4/8] fix: logging num_output_rows_total none return none instead of 1 Signed-off-by: machichima --- python/ray/data/_internal/progress/logging_progress.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/python/ray/data/_internal/progress/logging_progress.py b/python/ray/data/_internal/progress/logging_progress.py index 6530c7202315..5fabfb859bdb 100644 --- a/python/ray/data/_internal/progress/logging_progress.py +++ b/python/ray/data/_internal/progress/logging_progress.py @@ -124,7 +124,7 @@ def __init__( op = state.op if isinstance(op, InputDataBuffer): continue - total = op.num_output_rows_total() or 1 + total = op.num_output_rows_total() contains_sub_progress_bars = isinstance(op, SubProgressBarMixin) sub_progress_bar_enabled = show_op_progress and ( From 535e2f2fba69c12e46429598f9a0f84ac78ab603 Mon Sep 17 00:00:00 2001 From: machichima Date: Sun, 27 Sep 2026 10:37:43 +0800 Subject: [PATCH 5/8] feat: non aggre return total row from input Signed-off-by: machichima --- .../shuffle_operators/external_shuffle_map_operator.py | 4 +--- .../operators/shuffle_operators/shuffle_map_operator.py | 4 +--- 2 files changed, 2 insertions(+), 6 deletions(-) diff --git a/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_map_operator.py b/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_map_operator.py index b5daa2cee1bc..64ad629ed49d 100644 --- a/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_map_operator.py +++ b/python/ray/data/_internal/execution/operators/shuffle_operators/external_shuffle_map_operator.py @@ -108,7 +108,6 @@ def __init__( self._partition_bundles_emitted: bool = False # -- Stats ----------------------------------------------------------- - self._total_input_rows: int = 0 self._total_input_bytes: int = 0 self._map_blocks_stats: List[BlockStats] = [] # Per-partition decoded stats summed across completed mappers: @@ -327,7 +326,6 @@ def _handle_map_done( for bundle in input_bundles: bundle.destroy_if_owned() - self._total_input_rows += input_rows self._total_input_bytes += input_bytes input_meta = BlockMetadata( num_rows=input_rows, @@ -518,7 +516,7 @@ def num_output_rows_total(self) -> Optional[int]: # until the maps run. if self._block_transformer is not None: return None - return self._total_input_rows if self._total_input_rows > 0 else None + return self.input_dependencies[0].num_output_rows_total() def current_logical_usage(self) -> ExecutionResources: return ExecutionResources( diff --git a/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_map_operator.py b/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_map_operator.py index f44380bcba2a..643c52ed6523 100644 --- a/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_map_operator.py +++ b/python/ray/data/_internal/execution/operators/shuffle_operators/shuffle_map_operator.py @@ -139,7 +139,6 @@ def __init__( self._partition_bundles_emitted: bool = False # -- Stats ----------------------------------------------------------- - self._total_input_rows: int = 0 self._total_input_bytes: int = 0 self._map_blocks_stats: List[BlockStats] = [] @@ -293,7 +292,6 @@ def _handle_map_done( for bundle in input_bundles: bundle.destroy_if_owned() - self._total_input_rows += input_meta.num_rows or 0 self._total_input_bytes += input_meta.size_bytes or 0 self._map_blocks_stats.append(input_meta.to_stats()) @@ -398,7 +396,7 @@ def num_output_rows_total(self) -> Optional[int]: # until the maps run. if self._block_transformer is not None: return None - return self._total_input_rows if self._total_input_rows > 0 else None + return self.input_dependencies[0].num_output_rows_total() def current_logical_usage(self) -> ExecutionResources: return ExecutionResources( From b9d81a4fe6c0213c0d5e78db9f1684915d202a80 Mon Sep 17 00:00:00 2001 From: machichima Date: Sun, 27 Sep 2026 10:37:59 +0800 Subject: [PATCH 6/8] test: add tests Signed-off-by: machichima --- python/ray/data/tests/test_hash_shuffle_v2.py | 57 +++++++++++++++++++ 1 file changed, 57 insertions(+) diff --git a/python/ray/data/tests/test_hash_shuffle_v2.py b/python/ray/data/tests/test_hash_shuffle_v2.py index ded44b0f3fd4..59f3f752cc34 100644 --- a/python/ray/data/tests/test_hash_shuffle_v2.py +++ b/python/ray/data/tests/test_hash_shuffle_v2.py @@ -9,6 +9,9 @@ from ray.data._internal.execution.operators.shuffle_operators.external_shuffle_map_operator import ( # noqa: E501 ExternalHashShuffleMapOp, ) +from ray.data._internal.execution.operators.shuffle_operators.external_shuffle_reduce_operator import ( # noqa: E501 + ExternalHashShuffleReduceOp, +) from ray.data._internal.execution.operators.shuffle_operators.shuffle_map_operator import ( # noqa: E501 ShuffleMapOp, make_partition_sentinel, @@ -435,6 +438,60 @@ def test_reduce_op_runs_when_an_input_is_missing(ray_start_regular_shared_2_cpus assert op.has_completed() +_V2_OP_CLASSES = [ + (ShuffleMapOp, ShuffleReduceOp), + (ExternalHashShuffleMapOp, ExternalHashShuffleReduceOp), +] + + +def _make_map_op(map_op_cls, upstream_total_rows=None, block_transformer=None): + ctx = DataContext.get_current() + upstream = InputDataBuffer(ctx, []) + upstream.num_output_rows_total = lambda: upstream_total_rows + return map_op_cls( + upstream, + ctx, + num_partitions=2, + partition_fn=lambda table: {}, + block_transformer=block_transformer, + ) + + +@pytest.mark.parametrize("map_op_cls", [ShuffleMapOp, ExternalHashShuffleMapOp]) +def test_map_num_output_rows_total_unknown_with_block_transformer(map_op_cls): + map_op = _make_map_op( + map_op_cls, upstream_total_rows=100, block_transformer=lambda table: table + ) + + assert map_op.num_output_rows_total() is None + + +@pytest.mark.parametrize("map_op_cls,reduce_op_cls", _V2_OP_CLASSES) +def test_reduce_num_output_rows_total_borrows_map_total(map_op_cls, reduce_op_cls): + map_op = _make_map_op(map_op_cls, upstream_total_rows=100) + reduce_op = reduce_op_cls( + map_op, map_op.data_context, num_partitions=2, reduce_fn=lambda *args: [] + ) + + assert reduce_op.num_output_rows_total() == 100 + + +@pytest.mark.parametrize("map_op_cls,reduce_op_cls", _V2_OP_CLASSES) +def test_reduce_num_output_rows_total_unknown_when_row_count_not_preserved( + map_op_cls, reduce_op_cls +): + map_op = _make_map_op(map_op_cls, upstream_total_rows=100) + reduce_op = reduce_op_cls( + map_op, + map_op.data_context, + num_partitions=2, + reduce_fn=lambda *args: [], + preserves_row_count=False, + ) + + assert reduce_op.num_output_rows_total() is None + + if __name__ == "__main__": import sys From f2391107ae6d7001c44cbdfbd001182671fc4e77 Mon Sep 17 00:00:00 2001 From: machichima Date: Sun, 27 Sep 2026 16:13:10 +0800 Subject: [PATCH 7/8] fix: format for 0 show 0 rather than ? Signed-off-by: machichima --- python/ray/data/_internal/progress/logging_progress.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/python/ray/data/_internal/progress/logging_progress.py b/python/ray/data/_internal/progress/logging_progress.py index 5fabfb859bdb..4bff004e0607 100644 --- a/python/ray/data/_internal/progress/logging_progress.py +++ b/python/ray/data/_internal/progress/logging_progress.py @@ -218,7 +218,8 @@ def update_operator_progress( def _format_progress(m: _LoggingMetrics) -> str: - return f"{m.name}: {m.completed}/{m.total or '?'}" + total = "?" if m.total is None else m.total + return f"{m.name}: {m.completed}/{total}" def _log_global_progress(m: _LoggingMetrics): From 45a288cf0143fdd6ccec118732cb28b63f1c732c Mon Sep 17 00:00:00 2001 From: machichima Date: Sun, 27 Sep 2026 16:51:57 +0800 Subject: [PATCH 8/8] fix: rich view completed value source Signed-off-by: machichima --- python/ray/data/_internal/progress/rich_progress.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/python/ray/data/_internal/progress/rich_progress.py b/python/ray/data/_internal/progress/rich_progress.py index d66fac911ab5..7aec1bcdddda 100644 --- a/python/ray/data/_internal/progress/rich_progress.py +++ b/python/ray/data/_internal/progress/rich_progress.py @@ -280,8 +280,8 @@ def close_with_finishing_description(self, desc: str, success: bool): pg.complete() if self._start_time is None: self._start_time = time.time() - for tid, progress, _ in self._op_display.values(): - completed = progress.tasks[tid].completed or 0 + for op_state, (tid, progress, _) in self._op_display.items(): + completed = op_state.op.metrics.row_outputs_taken metrics = _get_progress_metrics( self._start_time, completed, completed )