API ReferenceExecutors

XennaExecutor

View as Markdown

XennaExecutor is the production executor that uses Cosmos-Xenna for distributed execution. It’s the default executor used when running pipelines.

Import

from nemo_curator.backends.xenna import XennaExecutor

Class Definition

class XennaExecutor(BaseExecutor):
"""Production executor using Cosmos-Xenna for distributed execution.
Provides:
- Distributed task orchestration
- Resource allocation and management
- Batch processing optimization
- Performance metrics collection
"""
def __init__(
self,
config: dict[str, Any] | None = None,
ignore_head_node: bool = False,
) -> None:
"""Initialize the executor.
Args:
config: Executor configuration dictionary.
ignore_head_node: Not supported. Raises ValueError if True.
"""

Configuration Options

OptionTypeDefaultDescription
logging_intervalint60Seconds between progress logs
ignore_failuresboolFalseContinue on task failures
execution_modestr"streaming""streaming" or "batch"
cpu_allocation_percentagefloat0.95CPU allocation fraction
autoscale_interval_sint180Autoscaling check interval

Usage Examples

Default Configuration

from nemo_curator.pipeline import Pipeline
from nemo_curator.backends.xenna import XennaExecutor
pipeline = Pipeline(name="my_pipeline", stages=[...])
# Default executor
results = pipeline.run()
# Equivalent to:
executor = XennaExecutor()
results = pipeline.run(executor=executor)

Per-Stage Worker Overrides

Use num_workers for a cluster-wide count. To request workers per node instead, leave num_workers unset and configure only num_workers_per_node:

from nemo_curator.stages.base import ProcessingStage
from nemo_curator.stages.resources import Resources
from nemo_curator.tasks import EmptyTask
class PassthroughStage(ProcessingStage[EmptyTask, EmptyTask]):
name = "passthrough"
resources = Resources(cpus=1.0)
def process(self, task: EmptyTask) -> EmptyTask:
return task
cluster_wide_stage = PassthroughStage().with_(
num_workers=12,
xenna_stage_spec={
"worker_max_lifetime_m": 45,
},
)
per_node_stage = PassthroughStage().with_(
num_workers=None,
xenna_stage_spec={
"num_workers_per_node": 2,
},
)

Do not set both worker controls. Do not place num_workers inside xenna_stage_spec; use with_(num_workers=...) or override the stage’s num_workers() method.

See Stage Worker Sizing for scope, invalid-combination behavior, and with_() merge semantics.

Custom Configuration

executor = XennaExecutor(config={
"logging_interval": 30,
"ignore_failures": True,
"execution_mode": "batch",
"cpu_allocation_percentage": 0.9,
})
results = pipeline.run(executor=executor)

Streaming vs Batch Mode

Processes tasks as they become available:

executor = XennaExecutor(config={
"execution_mode": "streaming",
})

Best for:

  • Large datasets
  • Memory-constrained environments
  • Real-time processing

Methods

execute()

Execute the pipeline stages.

def execute(
self,
stages: list[ProcessingStage],
initial_tasks: list[Task] | None = None,
) -> list[Task]:
"""Execute the pipeline stages.
Args:
stages: List of processing stages to execute.
initial_tasks: Initial tasks (defaults to a fresh `EmptyTask()` instance).
Returns:
List of output tasks from the final stage.
"""

Error Handling

executor = XennaExecutor(config={
"ignore_failures": True, # Continue despite errors
})
try:
results = pipeline.run(executor=executor)
except Exception as e:
# Handle pipeline-level failures
print(f"Pipeline failed: {e}")

Performance Monitoring

The executor automatically collects performance metrics:

results = pipeline.run(executor=executor)
# Each task contains performance data
for task in results:
for perf in task._stage_perf:
print(f"Stage: {perf.stage_name}")
print(f" Duration: {perf.process_time}s")
print(f" Items processed: {perf.num_items_processed}")

Source Code

View source on GitHub