nemo_curator.backends.ray_data.diagnostics

View as Markdown

Runtime Ray Data scheduler diagnostics for the supported Ray release.

The diagnostics were developed as a Ray source patch, but Curator cannot distribute a patched Ray wheel. This module installs the equivalent Python hooks in the driver process. The hooks emit through child loggers of ray.data so Ray’s SessionFileHandler writes them to session_latest/logs/ray-data/ray-data.log.

All affected scheduler components run in the Ray Data driver. Worker environments therefore do not need modified Ray installations.

Module Contents

Classes

NameDescription
DiagnosticsInstallStatusResult of attempting to enable Ray Data diagnostics.
_TaskAdmissionDecision-

Functions

NameDescription
_has_native_diagnostics-
_install_actor_autoscaling_diagnostics-
_install_downstream_capacity_diagnostics-
_install_resource_admission_diagnostics-
_install_scheduling_reasons-
_milliseconds-
_object_store_memory_fields-
execution_resource_fieldsFlatten Ray ExecutionResources into stable scalar fields.
format_logfmt_eventFormat an event as stable, parseable logfmt-like tokens.
install_ray_data_diagnosticsInstall driver-side diagnostics without modifying the Ray installation.

Data

RAY_DATA_DIAGNOSTICS_ENV_VAR

_INSTALL_LOCK

_INSTALL_MARKER

_SUPPORTED_RAY_VERSION

_TRUE_ENV_VALUES

API

class nemo_curator.backends.ray_data.diagnostics.DiagnosticsInstallStatus

Bases: enum.Enum

Result of attempting to enable Ray Data diagnostics.

ALREADY_INSTALLED
= 'already_installed'
DISABLED
= 'disabled'
INSTALLED
= 'installed'
NATIVE
= 'native'
UNSUPPORTED
= 'unsupported'
class nemo_curator.backends.ray_data.diagnostics._TaskAdmissionDecision(
allowed: bool,
reason: str,
incremental_resources: typing.Any,
remaining_budget: typing.Any,
pending_output_estimate: float | None,
op_usage: typing.Any = None,
allocation: typing.Any = None
)
Dataclass
allowed
bool
pending_output_estimate
float | None
reason
str
nemo_curator.backends.ray_data.diagnostics._has_native_diagnostics(
autoscaler_module: typing.Any,
resource_manager_module: typing.Any,
executor_state_module: typing.Any
) -> bool
nemo_curator.backends.ray_data.diagnostics._install_actor_autoscaling_diagnostics(
autoscaler_module: typing.Any
) -> None
nemo_curator.backends.ray_data.diagnostics._install_downstream_capacity_diagnostics(
downstream_policy_module: typing.Any
) -> None
nemo_curator.backends.ray_data.diagnostics._install_resource_admission_diagnostics(
resource_manager_module: typing.Any,
resource_policy_module: typing.Any
) -> None
nemo_curator.backends.ray_data.diagnostics._install_scheduling_reasons(
executor_state_module: typing.Any
) -> None
nemo_curator.backends.ray_data.diagnostics._milliseconds(
seconds: float
) -> float
nemo_curator.backends.ray_data.diagnostics._object_store_memory_fields(
resource_manager: typing.Any,
op: typing.Any
) -> dict[str, object]
nemo_curator.backends.ray_data.diagnostics.execution_resource_fields(
prefix: str,
resources: typing.Any
) -> dict[str, object]

Flatten Ray ExecutionResources into stable scalar fields.

nemo_curator.backends.ray_data.diagnostics.format_logfmt_event(
event: str,
fields: dict[str, object]
) -> str

Format an event as stable, parseable logfmt-like tokens.

nemo_curator.backends.ray_data.diagnostics.install_ray_data_diagnostics() -> nemo_curator.backends.ray_data.diagnostics.DiagnosticsInstallStatus

Install driver-side diagnostics without modifying the Ray installation.

Diagnostics are opt-in through NEMO_CURATOR_RAY_DATA_DIAGNOSTICS. The shim is intentionally restricted to the Ray version whose private APIs it targets. A future Ray release containing the upstream diagnostics is detected and left untouched.

nemo_curator.backends.ray_data.diagnostics.RAY_DATA_DIAGNOSTICS_ENV_VAR = 'NEMO_CURATOR_RAY_DATA_DIAGNOSTICS'
nemo_curator.backends.ray_data.diagnostics._INSTALL_LOCK = threading.Lock()
nemo_curator.backends.ray_data.diagnostics._INSTALL_MARKER = '_nemo_curator_ray_data_diagnostics_installed'
nemo_curator.backends.ray_data.diagnostics._SUPPORTED_RAY_VERSION = '2.57.0'
nemo_curator.backends.ray_data.diagnostics._TRUE_ENV_VALUES = {'1', 'true', 'yes', 'on'}