Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
Expand All @@ -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__)


Expand All @@ -58,9 +54,7 @@ def _make_mapper_sentinel(mapper_id: int) -> Tuple[str, ...]:
return (f"{_MAPPER_ID_SENTINEL}{mapper_id}",)


class DiskHashShuffleMapOp(
InternalQueueOperatorMixin, PhysicalOperator, SubProgressBarMixin
):
class DiskHashShuffleMapOp(InternalQueueOperatorMixin, PhysicalOperator):
"""Disk-shuffle map operator. See module docstring."""

_DEFAULT_SHUFFLE_MAP_TASK_NUM_CPUS = 1.0
Expand Down Expand Up @@ -114,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:
Expand All @@ -124,9 +117,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

# =====================================================================
# Disk-shuffle-specific state below.
# =====================================================================
Expand Down Expand Up @@ -274,15 +264,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,
Expand Down Expand Up @@ -345,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,
Expand All @@ -365,9 +345,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:
Expand Down Expand Up @@ -534,7 +511,12 @@ def get_stats(self) -> Dict[str, List[BlockStats]]:
return {self._name: self._map_blocks_stats}

def num_output_rows_total(self) -> Optional[int]:
return self._total_input_rows if self._total_input_rows > 0 else None
# 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.input_dependencies[0].num_output_rows_total()

def current_logical_usage(self) -> ExecutionResources:
return ExecutionResources(
Expand Down Expand Up @@ -567,13 +549,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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,20 +39,18 @@
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

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__)


class DiskHashShuffleReduceOp(PhysicalOperator, SubProgressBarMixin):
class DiskHashShuffleReduceOp(PhysicalOperator):
"""Disk-shuffle reduce operator.

Structurally mirrors ``ShuffleReduceOp``: one wrapper bundle per partition
Expand Down Expand Up @@ -81,6 +79,7 @@ def __init__(
peak_memory_multiplier: float = SHUFFLE_PEAK_MEMORY_MULTIPLIER,
name: str = "DiskHashShuffleReduce",
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,
Expand All @@ -103,6 +102,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 -----------------------------------
Expand Down Expand Up @@ -130,9 +132,6 @@ def __init__(
# -- Stats -----------------------------------------------------------
self._output_blocks_stats: List[BlockStats] = []

# -- Sub-progress bars -----------------------------------------------
self._reduce_bar: Optional["BaseProgressBar"] = None

# =====================================================================
# Disk-shuffle-specific state below.
# =====================================================================
Expand Down Expand Up @@ -316,8 +315,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
Expand All @@ -343,11 +340,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,
Expand Down Expand Up @@ -400,10 +392,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, DiskHashShuffleMapOp)
Expand Down Expand Up @@ -439,10 +431,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
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
import dataclasses
import functools
import logging
import typing
from collections import defaultdict
from typing import Any, Dict, List, Optional, Tuple

Expand All @@ -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,
Expand All @@ -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__)


Expand All @@ -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):

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

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)
Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit 7259ee4. Configure here.

"""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
Expand Down Expand Up @@ -145,13 +139,9 @@ 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] = []

# -- Sub-progress bars -----------------------------------------------
self._map_bar: Optional["BaseProgressBar"] = None

@property
def _input_queues(self) -> List[BaseBundleQueue]:
return []
Expand Down Expand Up @@ -268,15 +258,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,
Expand Down Expand Up @@ -311,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())

Expand All @@ -322,9 +302,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:
Expand Down Expand Up @@ -414,7 +391,12 @@ def get_stats(self) -> Dict[str, List[BlockStats]]:
return {self._name: self._map_blocks_stats}

def num_output_rows_total(self) -> Optional[int]:
return self._total_input_rows if self._total_input_rows > 0 else None
# 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.input_dependencies[0].num_output_rows_total()

def current_logical_usage(self) -> ExecutionResources:
return ExecutionResources(
Expand Down Expand Up @@ -447,10 +429,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
Loading
Loading