zephon.observability
Observability configuration and helpers used across Zephon.
-
class ExecutionTrackingMode(*values)[source]
Bases: str, Enum
Controls how much execution detail is collected.
-
OFF = 'off'
-
STAGES = 'stages'
-
NODES = 'nodes'
-
classmethod from_value(value)[source]
- Return type:
ExecutionTrackingMode
-
property collects_nodes: bool
-
property collects_stages: bool
-
class FetchStageSummary(index, totals=<factory>, shard_totals=<factory>)[source]
Bases: object
Stage fetch timings keyed by (dataset_id, shard_id).
-
index: int
-
totals: FetchTimingTotals
-
shard_totals: MutableMapping[tuple[int, int], FetchTimingTotals]
-
apply(delta)[source]
-
merge(other)[source]
-
copy()[source]
- Return type:
FetchStageSummary
-
to_records(*, plan_id, stage_name, tracking_mode, include_shards=True)[source]
- Return type:
list[dict[str, Any]]
-
class FetchTimingDelta(stage_index, dataset_id, shard_id, samples, group_ns, resolve_ns, open_ns, read_ns, close_ns, retries, cache_hits, cache_misses, shard_reopens=0)[source]
Bases: object
Incremental fetch metrics emitted per shard group.
-
stage_index: int
-
dataset_id: int
-
shard_id: int
-
samples: int
-
group_ns: int
-
resolve_ns: int
-
open_ns: int
-
read_ns: int
-
close_ns: int
-
retries: int
-
cache_hits: int
-
cache_misses: int
-
shard_reopens: int
-
class FetchTimingSummary(plan_id=None, reporting_interval_s=None, tracking_mode=ExecutionTrackingMode.OFF, stages=<factory>)[source]
Bases: object
Aggregated fetch metrics for an entire pipeline.
-
plan_id: str | None = None
-
reporting_interval_s: float | None = None
-
tracking_mode: ExecutionTrackingMode = 'off'
-
stages: MutableMapping[int, FetchStageSummary]
-
apply(delta)[source]
-
merge(other)[source]
-
iter_stages()[source]
- Return type:
Iterator[FetchStageSummary]
-
to_records(stage_names=None, *, include_shards=True)[source]
- Return type:
list[dict[str, Any]]
-
clone()[source]
- Return type:
FetchTimingSummary
-
has_samples()[source]
- Return type:
bool
-
class FetchTimingTotals(samples=0, groups=0, group_ns=0, resolve_ns=0, open_ns=0, read_ns=0, close_ns=0, retries=0, cache_hits=0, cache_misses=0, shard_reopens=0)[source]
Bases: object
Aggregated timing totals accumulated across deltas.
-
samples: int
-
groups: int
-
group_ns: int
-
resolve_ns: int
-
open_ns: int
-
read_ns: int
-
close_ns: int
-
retries: int
-
cache_hits: int
-
cache_misses: int
-
shard_reopens: int
-
apply(delta)[source]
-
merge(other)[source]
-
copy()[source]
- Return type:
FetchTimingTotals
-
property avg_group_ns: float
-
property avg_resolve_ns: float
-
property avg_open_ns: float
-
property avg_read_ns: float
-
property avg_close_ns: float
-
property cache_hit_ratio: float
-
class MTPQueueStats(depth, capacity, staged_bytes, prefetch_depth)[source]
Bases: object
Point-in-time occupancy of the MTP hand-off queue (data_q).
depth counts queued items, which may not have been pickled yet;
staged_bytes is the already-pickled data in the kernel transport
buffer, readable without the feeder running. Zero staged with depth
near capacity means the feeder is starved. prefetch_depth counts
items already drained into the main-process prefetch buffer; it is
sampled without locking and may be slightly stale.
A field is -1 when unmeasurable: depth on macOS (sem_getvalue is
unsupported there), staged_bytes once the queue is closed.
-
depth: int
Items put by the subprocess but not yet retrieved by the consumer.
-
capacity: int
Maximum depth — the resolved mtp_buffer.
-
staged_bytes: int
Pickled bytes staged in the transport buffer, consumer-ready.
-
prefetch_depth: int
Deserialized items in the prefetch buffer (0 when mtp_prefetch=0).
-
class MetricsSinkConfig(mode=MetricsSinkMode.LOG, namespace='zephon.pipeline', statsd_host='127.0.0.1', statsd_port=8125, default_tags=<factory>, json_logs=True, max_batch_size=32, flush_interval_s=5.0)[source]
Bases: object
Describe how observability data is exported.
-
mode: MetricsSinkMode = 'log'
-
namespace: str = 'zephon.pipeline'
-
statsd_host: str = '127.0.0.1'
-
statsd_port: int = 8125
-
default_tags: Mapping[str, str]
-
json_logs: bool = True
-
max_batch_size: int = 32
-
flush_interval_s: float = 5.0
-
should_log()[source]
- Return type:
bool
-
should_emit_statsd()[source]
- Return type:
bool
-
with_additional_tags(extra)[source]
- Return type:
MetricsSinkConfig
-
class MetricsSinkMode(*values)[source]
Bases: str, Enum
Where metrics should be emitted.
-
LOG = 'log'
-
STATSD = 'statsd'
-
BOTH = 'both'
-
classmethod from_value(value)[source]
- Return type:
MetricsSinkMode
-
includes_logs()[source]
- Return type:
bool
-
includes_statsd()[source]
- Return type:
bool
-
class PipelineSummary(plan_id=None, reporting_interval_s=None, tracking_mode=ExecutionTrackingMode.OFF, stages=<factory>, total_wait_ratio=None)[source]
Bases: object
Aggregated metrics for an entire pipeline invocation.
-
plan_id: str | None = None
-
reporting_interval_s: float | None = None
-
tracking_mode: ExecutionTrackingMode = 'off'
-
stages: MutableMapping[int, StageSummary]
-
total_wait_ratio: float | None = None
-
apply(delta)[source]
-
merge(other)[source]
-
iter_stages()[source]
- Return type:
Iterator[StageSummary]
-
to_records()[source]
- Return type:
list[dict[str, Any]]
-
clone()[source]
- Return type:
PipelineSummary
-
compute_wait_ratios()[source]
Derive wait ratios for the pipeline based on processing and wait times.
-
class PrefetchTimingDelta(stage_index, batch_size, prefetch_requests, prefetch_succeeded, prefetch_failed)[source]
Bases: object
Incremental prefetch metrics emitted per batch.
-
stage_index: int
-
batch_size: int
-
prefetch_requests: int
-
prefetch_succeeded: int
-
prefetch_failed: int
-
class PrefetchTimingSummary(plan_id=None, reporting_interval_s=None, tracking_mode=ExecutionTrackingMode.OFF, stages=<factory>)[source]
Bases: object
Aggregated prefetch metrics for an entire pipeline.
-
plan_id: str | None = None
-
reporting_interval_s: float | None = None
-
tracking_mode: ExecutionTrackingMode = 'off'
-
stages: MutableMapping[int, PrefetchTimingTotals]
-
apply(delta)[source]
-
merge(other)[source]
-
clone()[source]
- Return type:
PrefetchTimingSummary
-
has_samples()[source]
- Return type:
bool
-
property success_rate: float
-
class PrefetchTimingTotals(batches=0, samples=0, prefetch_requests=0, prefetch_succeeded=0, prefetch_failed=0)[source]
Bases: object
Aggregated prefetch totals accumulated across deltas.
-
batches: int
-
samples: int
-
prefetch_requests: int
-
prefetch_succeeded: int
-
prefetch_failed: int
-
apply(delta)[source]
-
merge(other)[source]
-
copy()[source]
- Return type:
PrefetchTimingTotals