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 exceed MAX_SHARD_FAILURES (default 3).
  • INFRA-type: Preemption or heartbeat timeouts counted in task_infra_attempts. These abort after MAX_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 > 0 with queue_depth == 0 in ZephyrCoordinator.get_status and high resource usage in per-worker metrics.
  • Debug coordinator stalls by examining _coordinator_loop logs for worker death patterns and monitoring progress_time_seconds telemetry.
  • 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_worker and _maybe_requeue_worker_tasks, but investigate if re-registration occurs frequently.
  • Profile bottlenecks using the Iris CLI profile-task command 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →