Finelog Stats Namespaces for Task and Worker Telemetry in Marin: A Complete Reference

Marin uses two dedicated Finelog namespaces—zephyr.stage for pipeline stage aggregates and zephyr.worker for per-shard worker metrics—to capture telemetry from distributed Zephyr executions.

Marin pipelines rely on Finelog as the primary telemetry sink for monitoring distributed training and inference workloads. The Zephyr stats module defines dedicated string constants that partition metrics into logical streams: one for high-level stage completion statistics and another for granular worker heartbeat data. Understanding these Finelog stats namespaces is essential for querying execution history and debugging pipeline performance in the Marin ecosystem.

Understanding the Two Finelog Namespaces

Marin separates telemetry into two distinct namespaces to optimize query patterns and storage cardinality.

Stage Statistics Namespace (zephyr.stage)

The ZEPHYR_STAGE_STATS_NAMESPACE constant maps to the string "zephyr.stage" as defined in lib/zephyr/src/zephyr/stats.py at line 34. This namespace receives exactly one structured row per pipeline stage upon completion or failure, containing aggregated counters such as total items processed, bytes transferred, and average CPU utilization across all shards.

Worker Statistics Namespace (zephyr.worker)

The ZEPHYR_WORKER_STATS_NAMESPACE constant equals "zephyr.worker" (line 35 of lib/zephyr/src/zephyr/stats.py). Unlike the stage namespace, the worker namespace receives multiple rows per execution: one when a shard starts, periodic samples at configurable intervals, and a final row when the shard terminates. This enables real-time monitoring of resource consumption at the individual worker level.

Namespace Implementation in Zephyr Stats

Both constants are implemented as module-level strings in the Zephyr stats library:


# lib/zephyr/src/zephyr/stats.py (lines 34-35)

ZEPHYR_STAGE_STATS_NAMESPACE = "zephyr.stage"
ZEPHYR_WORKER_STATS_NAMESPACE = "zephyr.worker"

These values are consumed by the StatsWriter class, which wraps a Finelog LogClient to handle serialization and transport. The StatsWriter abstracts the underlying Finelog connection, automatically resolving endpoints via Iris when available.

Emitting Telemetry Using StatsWriter

Applications interact with these namespaces through the StatsWriter helper rather than writing directly to Finelog. The class provides two primary methods: emit_stage_stat() for stage-level aggregates and emit_worker_stat() for shard-level telemetry.

Writing Stage Statistics

After a pipeline stage completes, use emit_stage_stat() to write aggregated metrics to the zephyr.stage namespace:

from zephyr.stats import (
    ZEPHYR_STAGE_STATS_NAMESPACE,
    StatsWriter,
    ZephyrWorkerStatStatus,
)

stats = StatsWriter.connect()

# Aggregate counters from all shards

stage_counters = {
    "zephyr/item_count": 1_200,
    "zephyr/bytes_processed": 45_000_000,
    "zephyr/worker/cpu_pct_average": 78.3,
}

stats.emit_stage_stat(
    stage_counters=stage_counters,
    stage_name="training",
    execution_id="run-2024-08-29-01",
    elapsed=360.5,
    total_shards=16,
    status=ZephyrWorkerStatStatus.END,
)

Writing Worker Statistics

For per-shard telemetry, including periodic heartbeats, use emit_worker_stat() targeting the zephyr.worker namespace:

from datetime import datetime

stats.emit_worker_stat(
    shard_idx=3,
    stage_name="training",
    execution_id="run-2024-08-29-01",
    status=ZephyrWorkerStatStatus.RUNNING,
    ts=datetime.utcnow(),
    items=300,
    bytes_processed=12_500_000,
    item_rate=0.83,
    byte_rate=34_500.0,
    cpu_time_total=120.0,
    cpu_current_pct=79.1,
    cpu_avg_pct=78.5,
    mem_current_bytes=2_048_000_000,
    mem_avg_bytes=2_030_000_000,
    mem_peak_bytes=2_100_000_000,
)

Key Source Files and Architecture

The telemetry system spans three primary components in the Marin repository:

Summary

  • Marin pipelines use two Finelog stats namespaces: zephyr.stage for aggregated stage metrics and zephyr.worker for per-shard telemetry.
  • Constants are defined in lib/zephyr/src/zephyr/stats.py as ZEPHYR_STAGE_STATS_NAMESPACE (line 34) and ZEPHYR_WORKER_STATS_NAMESPACE (line 35).
  • The StatsWriter class provides emit_stage_stat() and emit_worker_stat() methods to write to these namespaces.
  • Stage entries are written once per stage completion, while worker entries capture lifecycle events and periodic samples for each shard.

Frequently Asked Questions

How do I query stage-level telemetry in Finelog?

Query the zephyr.stage namespace and filter by execution_id and stage_name. Each row represents a completed stage execution with aggregated counters from all workers, emitted automatically when a stage finishes or fails.

What is the difference between stage and worker namespaces in Marin?

The stage namespace (zephyr.stage) receives a single row per pipeline stage containing aggregated statistics, while the worker namespace (zephyr.worker) receives multiple rows per execution—start events, periodic heartbeats, and end events—for each individual shard.

Where are the Finelog namespace constants defined in the Marin source code?

The constants ZEPHYR_STAGE_STATS_NAMESPACE and ZEPHYR_WORKER_STATS_NAMESPACE are defined in lib/zephyr/src/zephyr/stats.py at lines 34 and 35, respectively.

Can I write custom metrics to these namespaces?

Yes. Use the StatsWriter class methods with custom counter keys defined in your application code or added to lib/zephyr/src/zephyr/counters.py. Ensure counter keys follow the zephyr/ prefix convention for consistency with built-in metrics.

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 →