# How to Debug Zephyr Pipeline Stragglers and Coordinator Issues in marin-community/marin

> Debug Zephyr pipeline stragglers and coordinator issues in marin-community/marin. Monitor counters and heartbeat loss to identify and resolve problems quickly.

- Repository: [The Marin Project/marin](https://github.com/marin-community/marin)
- Tags: how-to-guide
- Published: 2026-08-29

---

**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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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:

```python
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:

```bash
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`](https://github.com/marin-community/marin/blob/main/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:

```bash
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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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:

```bash
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:

```bash
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:

```bash
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`):

```bash
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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/coordinator.py) contains the core implementation for `get_status`, `_coordinator_loop`, and failure handling logic.