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 |
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]
|
|
Dict[str, Any]
|
|
Dict[str, Any]
|
|
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 |
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]
|
|
Dict[str, Any]
|
|
cleanup()
¶
Clean up resources including workers and SHM pools.