How to Debug Zephyr Pipeline Stragglers and Coordinator Issues in marin-community/marin
Debug Zephyr pipeline stragglers by monitoring in_flight versus queue_depth counters in ZephyrCoordinator.get_status, while coordinator failure patterns surface via heartbeat loss in _coordinator_loop and automatic shard re-queuing in _maybe_requeue_worker_tasks.
Zephyr pipelines in the marin-community/marin repository implement a pull-based coordinator-worker architecture where the coordinator drives each stage and workers repeatedly call pull_task() to fetch shards. When pipelines stall or fail, the root cause typically falls into two categories: stragglers (slow or stuck shards) and coordinator issues (heartbeat loss, worker death, or stage stalls). Understanding the interaction between lib/zephyr/src/zephyr/coordinator.py and the worker pool is essential for effective debugging.
Identifying Pipeline Stragglers
Stragglers manifest as shards that execute significantly slower than their peers, causing the entire stage to wait. The coordinator tracks these via internal counters and per-worker resource metrics.
Detecting Stuck Shards via Status Counters
In ZephyrCoordinator.get_status at lib/zephyr/src/zephyr/coordinator.py (line 1010), the coordinator exposes in_flight, queue_depth, and completed counters. When in_flight >> 0 while queue_depth == 0, the pipeline is waiting on a few slow shards to finish. This state indicates that no new work is being assigned because the coordinator cannot proceed until the current batch completes.
To inspect this programmatically:
from zephyr.coordinator import ZephyrCoordinator
coord = ZephyrCoordinator(...) # injected by driver
status = coord.get_status()
if status.in_flight and status.queue_depth == 0:
print("Potential straggler detected: {} shards blocking stage".format(status.in_flight))
Analyzing Worker Resource Metrics
Resource-heavy shards often correlate with high memory or disk usage and low CPU utilization. Examine the workers array in the JobStatus JSON for entries with elevated memoryMb, diskMb, or low cpuPercent. Retrieve the full status via the Iris CLI:
uv run iris --config prod actor call <ENDPOINT> get_status
Sort the workers by resource consumption to identify outliers holding large shards.
Diagnosing Coordinator Issues
Coordinator problems range from heartbeat timeouts to complete worker pool failure. The _coordinator_loop function (line 63 in coordinator.py) drives the lifecycle of the pipeline and emits critical telemetry.
Monitoring the Coordinator Loop and Heartbeats
The coordinator loop logs a summary every approximately five seconds via _log_status. Look for the log line format:
[%s] [%s] %d/%d complete, %d in-flight, %d queued, %d/%d workers alive, %d dead
Abnormal worker counts in this line indicate instability. If all workers transition to FAILED or DONE states, the coordinator raises ZephyrWorkerError after no_workers_timeout (default 60 seconds). To retrieve these logs:
uv run iris rpc controller get-task-logs --id <JOB_ID> --max-total-lines 5000 --attempt-id -1 --tail
Search for "All workers are dead/failed" to identify total worker loss, which typically signals cluster-level resource exhaustion or pod OOMs.
Handling Worker Death and Re-registration
When a worker crashes and restarts, register_worker (line 141 in coordinator.py) updates the worker handle and invokes _maybe_requeue_worker_tasks to redistribute any in-flight shards. This can surface as sudden spikes in pull_task logs as the replacement worker claims the re-queued work. Monitor the task_error_attempts versus task_infra_attempts counters to ensure the shard is not repeatedly crashing, which would trigger a fatal pipeline abort.
Detecting Stage Stalls via Telemetry
Stalled stages are detected by the progress_time_seconds metric published in _publish_telemetry. The coordinator updates this metric on each shard completion. If Grafana reports ZephyrPipelineProgressStalled, call get_status to confirm that no shard completed within the last 45 minutes, indicating the coordinator is blocked on a straggler or dead workers.
Handling Shard Failures and Abort Conditions
The coordinator distinguishes between task logic failures and infrastructure failures, applying different retry thresholds before aborting.
TASK vs INFRA Failure Types
In _record_shard_failure (line 151 in coordinator.py), failures are categorized as:
- TASK-type: User-code errors counted in
task_error_attempts. The pipeline aborts when these exceedMAX_SHARD_FAILURES(default 3). - INFRA-type: Preemption or heartbeat timeouts counted in
task_infra_attempts. These abort afterMAX_SHARD_INFRA_FAILURES(default 20).
When a worker dies while holding a task, the coordinator automatically re-queues the shard as an INFRA failure via _maybe_requeue_worker_tasks. Excessive INFRA attempts suggest unstable workers rather than code bugs.
Step-by-Step Debugging Workflow
Follow this structured approach to resolve Zephyr pipeline issues.
1. Gather Coordinator Logs
Pull the latest logs and search for periodic status lines and shard failure warnings:
uv run iris rpc controller get-task-logs --id <JOB_ID> --max-total-lines 5000 --attempt-id -1 --tail
2. Query Real-Time Status
Retrieve the full JobStatus JSON to check worker health and queue state:
uv run iris --config <CONFIG> actor call <ENDPOINT> get_status
Compare workers entries to spot resource outliers.
3. Identify Straggler Shards
If in_flight is non-zero but queue_depth is zero, list tasks to find heavy shards:
uv run iris rpc controller list-tasks --job-id <JOB_ID> \
| jq '.tasks[] | {id, memoryMb, diskMb, cpuPercent}'
4. Profile the Suspect Worker
Generate a CPU profile to determine if the bottleneck is user-code (GIL contention) or I/O (write_table):
uv run iris rpc controller profile-task \
--json '{"target":"<WORKER_ID>","durationSeconds":10,"profileType":{"cpu":{"format":"SPEEDSCOPE"}}}'
5. Verify Telemetry and Custom Counters
Confirm that progress_time_seconds is being emitted in _publish_telemetry. For custom user counters merged in _aggregate_counter_snapshots, retrieve them via get_counters and compare across stages to detect anomalous throughput.
Summary
- Detect stragglers by checking
in_flight > 0withqueue_depth == 0inZephyrCoordinator.get_statusand high resource usage in per-worker metrics. - Debug coordinator stalls by examining
_coordinator_looplogs for worker death patterns and monitoringprogress_time_secondstelemetry. - Handle failures according to type: TASK failures abort after 3 attempts via
_record_shard_failure, while INFRA failures tolerate 20 attempts before fatal abort. - Recover workers automatically through
register_workerand_maybe_requeue_worker_tasks, but investigate if re-registration occurs frequently. - Profile bottlenecks using the Iris CLI
profile-taskcommand with Speedscope format to distinguish CPU-bound from I-bound operations.
Frequently Asked Questions
How do I distinguish between a straggler shard and a hung coordinator?
A straggler shard maintains in_flight > 0 with queue_depth == 0 in the status output, indicating active but slow work. A hung coordinator shows progress_time_seconds stagnation in Grafana with no recent shard completions, often accompanied by "All workers are dead/failed" log entries from _coordinator_loop.
What causes the coordinator to raise ZephyrWorkerError?
The coordinator raises ZephyrWorkerError when all workers enter FAILED or DONE states and remain unavailable for no_workers_timeout seconds (default 60). This typically indicates cluster-wide resource exhaustion, pod preemption, or cascading worker crashes rather than single-shard issues.
How does the coordinator handle a worker that dies mid-task?
When a worker crashes, the register_worker method updates the worker handle and _maybe_requeue_worker_tasks automatically re-queues any in-flight shards as INFRA failures. The replacement worker will pull these tasks via pull_task upon restart, though excessive re-queuing suggests underlying stability problems.
Where can I find operational playbooks for Zephyr debugging?
The lib/zephyr/OPS.md file contains the operational guide outlining dashboard commands, diagnostic patterns, and straggler detection workflows. Additionally, lib/zephyr/src/zephyr/coordinator.py contains the core implementation for get_status, _coordinator_loop, and failure handling logic.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →