Skip to content

Streaming orchestrator

cutana.streaming_orchestrator

Streaming Orchestrator for Cutana - handles batch-by-batch streaming workflows.

This module provides StreamingOrchestrator, a specialized orchestrator for processing large catalogues in batches with adaptive parallel processing.

Key features: - Adaptive parallelism: spawns workers based on consumer demand - Per-worker shared memory pools for zero-copy cutout transfer - Decoupled internal/user batch sizes to reduce subprocess overhead - Unordered result delivery: next_batch() returns whichever batch finishes first

StreamingOrchestrator(config, status_panel=None)

Bases: Orchestrator

Orchestrator specialized for streaming batch-by-batch processing.

Uses memory-efficient streaming with adaptive parallel workers. Each worker gets its own pre-allocated shared memory pool.

Usage

orchestrator = StreamingOrchestrator(config) try: orchestrator.init_streaming( batch_size=500, write_to_disk=False, max_workers=4, min_workers=2 )

for i in range(orchestrator.get_batch_count()):
    result = orchestrator.next_batch()
    # Process result["cutouts"] or result["zarr_path"]

finally: # Required: workers block indefinitely waiting to hand over their next # chunk, so abandoning the loop without cleanup() leaves them running # and holding their shared memory pools. orchestrator.cleanup()

Initialize the streaming orchestrator.

init_streaming(batch_size, write_to_disk=True, max_workers=4, min_workers=1, max_shm_memory_consumption=None)

Initialize streaming mode for batch-by-batch processing.

Parameters:

Name Type Description Default
batch_size int

Sources per user-facing batch. In-memory batches are cut to exactly this size from a rolling buffer, so it is an exact size there and is independent of how sources are split across worker processes. In disk mode each batch is one zarr archive written by one worker, so it cannot be re-split: there batch_size is a minimum, and :meth:_resolve_internal_batch_size may enlarge it.

required
write_to_disk bool

If True, write batches to zarr; if False, return cutouts in memory

True
max_workers int

Maximum number of parallel workers (default: 4)

4
min_workers int

Number of workers to pre-spawn at init time so they are already processing when next_batch() is first called (default: 1)

1
max_shm_memory_consumption Optional[int]

Total SHM memory budget in bytes (None = auto)

None

next_batch()

Get the next batch of cutouts.

Returns whichever worker finishes first (unordered). Adaptively spawns additional workers if the caller is consuming faster than production.

In-memory batches are assembled from a rolling buffer so that every batch except the last contains exactly batch_size cutouts; only the final batch may be smaller. Disk-mode batches map 1:1 to internal batches.

Returns:

Type Description
Dict[str, Any]

Dictionary with:

Dict[str, Any]
  • 'batch_number': 1-indexed batch number
Dict[str, Any]
  • 'cutouts': list of cutout arrays (if write_to_disk=False)
Dict[str, Any]
  • 'metadata': list of source metadata dicts

Raises:

Type Description
RuntimeError

If not initialized or no more batches

get_batch_count()

Get total number of user-facing batches.

get_worker_info()

Get per-worker detail keyed by process_id (issue #354).

Each :class:~cutana.profiling_types.WorkerInfo combines spawn-time batch composition (batch_index, n_sources, pool_slot, start_time) with completion-time fields (end_time, sources_per_fits_set, performance — the per-stage wall/cpu/stall_time/read_bytes breakdown). The completion-time fields are None until the worker reports completion via the in-memory streaming path. This supersedes the old worker-event list (it carries start/end times and more), so it is the single profiling entry point.

Returns:

Type Description
Dict[str, WorkerInfo]

A mapping of process_id to that worker's WorkerInfo.

get_delivery_report()

Report how many cutouts finished workers delivered against their assignment.

A source whose cutout window falls outside its tile yields nothing, so a run can legitimately produce fewer cutouts than the catalogue has rows. Such a deficit is only fatal when it empties a whole batch; otherwise the run completes and this is where the caller sees what was lost and which worker lost it.

Returns:

Type Description
Dict[str, Any]

Dictionary with:

Dict[str, Any]
  • 'missing': total cutouts assigned but never delivered so far
Dict[str, Any]
  • 'shortfalls': list of (process_id, assigned, delivered) tuples, one per worker that delivered fewer cutouts than it was assigned

cleanup()

Clean up resources including workers and SHM pools.