Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,30 @@ All notable changes are recorded here. Versions follow
The 0.1.0 entry describes what the library is rather than what changed, since
there was nothing to have changed from. Everything above it is ordinary.

## Unreleased

### Added

- **Metrics in the OTel example** (`examples/otel/metrics.py`):
- **`StreamMetricExporter`**, an OTel metric exporter that publishes each
collection to a stream, one row per data point. It covers sums, gauges
and both kinds of histogram, with exemplars kept, so a metric leads to
the trace it measured.
- **Delta temporality** for counters and histograms, so a window's total
is a `sum()` over the stored table.
- **The services record metrics:** an order count by outcome and
`http.server.request.duration`.
- **`otel/export.py` re-exports metrics as OTLP**, and **`just demo otel`
shows them in otel-gui's Metrics tab.** The dashboard moves to otel-gui
3.0.0, the first release that accepts metrics.

### Changed

- **The one-shot OTel demo runs without a maintainer**, as the migration demo
does. Its maintainer processes were still opening the logs when the run
ended, and cost it several seconds: `demo.main` took 3–11 s with them and
about 1 s without.

## 0.13.1 — 2026-10-03

### Changed
Expand Down
51 changes: 37 additions & 14 deletions examples/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ just demo NAME [ARGS] # run one, with all its processes
| `live` | as `trades`, with a live-only broker | no log: nothing to replay |
| `fastapi` | `fastapi_app.py` as the broker, then as `trades` | the broker mounted in a FastAPI app |
| `book` | broker, `book/producer.py`, `book/index.html` | Bitstamp's live order book as a keyed table log, kept by a browser page |
| `otel` | otel-gui, broker, `otel/services.py`, `otel/analytics.py`, `otel/export.py` | OpenTelemetry logs and traces through streams, in a dashboard |
| `otel` | otel-gui, broker, `otel/services.py`, `otel/analytics.py`, `otel/export.py` | OpenTelemetry logs, traces and metrics through streams, in a dashboard |
| `clean` | | delete what the demos stored |

`just demo` starts a demo's processes in order, waiting for each one that
Expand Down Expand Up @@ -260,22 +260,39 @@ against E. An empty diff means the migration does what production does. Once
it's empty, the cutover is C reading D's stream, and B retires with its
history still a queryable table.

## OpenTelemetry logs and traces, in a dashboard
## OpenTelemetry logs, traces and metrics, in a dashboard

```
just demo otel # otel-gui, a broker, the services, analytics, and the OTLP exporter
uv run python -m examples.otel.demo # the same pipeline once, printing what each part saw
```

Two simulated services, checkout and payments, handle traced orders and log
through Python's `logging` and the OpenTelemetry SDK. Their log records and
spans are published to two streams, `logs` and `spans`. `otel/export.py`
follows both and re-exports every row as OTLP to [otel-gui](https://github.com/metafab/otel-gui),
a local dashboard where logs, traces and the service map fill in live. A
Two simulated services, checkout and payments, handle traced orders, log
through Python's `logging`, and record metrics, all through the OpenTelemetry
SDK. Their log records, spans and metric data points are published to three
streams, `logs`, `spans` and `metrics`. `otel/export.py` follows all three
and re-exports every row as OTLP to [otel-gui](https://github.com/metafab/otel-gui),
a local dashboard where logs, traces, metrics and the service map fill in live. A
failed order is one trace across both services: payments' `POST /charge`
span with a `card declined` event and error log, and checkout's request
marked failed with an `order failed` warning.

**The metrics** are an order count by outcome (`shop.orders`) and each
service's request duration (`http.server.request.duration`, OTel's semantic
convention), exported every second. A stream row is one data point, carrying
its metric's name, unit, type and temporality. Counters and histograms are
exported as deltas, so each row is what happened in its own second and a
window's total is a `sum()`:

```sql
SELECT attributes['outcome'].string_value AS outcome, sum(value_int) AS orders
FROM log WHERE name = 'shop.orders' GROUP BY outcome
```

A measurement made inside a span keeps that span's trace as an *exemplar*,
so a metric leads straight to a request: the failed order's count points at
the failed order's trace.

`otel/analytics.py` is real-time analytics on the same telemetry, in one
`Stream.live` view of `spans`. Every few seconds it asks, in one SQL query,
each service's error rate over the last 30 seconds against its error rate
Expand All @@ -291,26 +308,29 @@ current, and the question is a query.

| file | what it holds |
|---|---|
| `otel/common.py` | what both signals share: `AnyValue`, ids, scope, and publishing from OTel's export thread |
| `otel/common.py` | what the signals share: `AnyValue`, ids, scope, and publishing from OTel's export thread |
| `otel/logs.py` | the log record's schema, its row conversion, and `StreamLogExporter` |
| `otel/spans.py` | the span's schema, its row conversion, and `StreamSpanExporter` |
| `otel/metrics.py` | the metric data point's schema, its row conversion, and `StreamMetricExporter` |
| `otel/demo.py` | the two services, the broker, and the one-shot demo |
| `otel/export.py` | rows back to OTel records and spans, out through OTel's OTLP exporters |
| `otel/services.py` | the producer: two services logging and tracing through OTel, publishing to the broker |
| `otel/export.py` | rows back to OTel records, spans and metric exports, out through OTel's OTLP exporters |
| `otel/services.py` | the producer: two services logging, tracing and recording metrics through OTel, publishing to the broker |
| `otel/analytics.py` | real-time analytics with `Stream.live`: error rate now against always, per service |
| `otel/gui.py` | otel-gui, downloaded, checked and run: where `export.py` sends what it reads |

**None of this is in streamcast.** The schemas and conversions are built from
the column types any stream can declare: trace and span ids as hex binary,
attributes as a map, OTel's `AnyValue` as a struct, and a span's events and
links as lists of structs. The OTel packages are dev dependencies, for this
attributes as a map, OTel's `AnyValue` as a struct, a span's events and
links as lists of structs, and a histogram's buckets as lists. The OTel packages are dev dependencies, for this
example only.

**`otel/export.py` is an ordinary OTLP exporter.** OTel viewers are
*receivers*: telemetry is pushed to them, and none subscribes to a WebSocket.
So this subscribes to the streams, turns each row back into an SDK log record
or span, and hands it to OpenTelemetry's own batch processors and OTLP/HTTP
exporters. The batching, the protobuf encoding and the retries are OTel's,
exporters. Metrics have no batch processor to hand a point to, so metric rows
are regrouped every half second into the resource, scope and metric nesting
an export carries, and handed to OTel's OTLP metric exporter. The batching, the protobuf encoding and the retries are OTel's,
and it works with any OTLP/HTTP receiver. Point `--receiver` at an OTel
Collector and it feeds whatever the Collector does.

Expand All @@ -324,4 +344,7 @@ The one-shot run shows what a subscriber can do beyond a dashboard:
- one failed request replayed by its trace id from both streams, with the id
given as hex text in `where=`;
- SQL over the stored tables: errors per service, the slowest request from
its root span, and ingest lag from `streamcast_ts`.
its root span, and ingest lag from `streamcast_ts`;
- SQL over the metrics: orders by outcome, mean request duration per service
from the histograms, and the failed order's exemplar, which is the failed
request's trace id.
4 changes: 3 additions & 1 deletion examples/__main__.py
Original file line number Diff line number Diff line change
Expand Up @@ -123,7 +123,7 @@ def trades(
],
),
"otel": Demo(
"OpenTelemetry logs and traces through streams, in otel-gui",
"OpenTelemetry logs, traces and metrics through streams, in otel-gui",
[
Process("otel-gui", module("examples.otel.gui"), 4318),
Process(
Expand All @@ -134,6 +134,8 @@ def trades(
"logs=examples.otel.logs:SCHEMA",
"--stream",
"spans=examples.otel.spans:SCHEMA",
"--stream",
"metrics=examples.otel.metrics:SCHEMA",
"--port",
"8766",
"--root",
Expand Down
81 changes: 67 additions & 14 deletions examples/otel/demo.py
Original file line number Diff line number Diff line change
@@ -1,28 +1,31 @@
"""OpenTelemetry logs and traces, through two streams, into tables you can query.
"""OpenTelemetry logs, traces and metrics, through streams, into tables you can query.

uv run python -m examples.otel.demo # once, printing what it saw
just demo otel # the same roles as processes, live

Nothing here is part of streamcast. The OTel schemas and conversions live in
`logs.py` and `spans.py`, built from ordinary column types. The OTel packages
`logs.py`, `spans.py` and `metrics.py`, built from ordinary column types. The OTel packages
are dev dependencies, for this example only.

The three roles, in one process so a test can run it: the broker, the
producer (`services.py`) and the subscribers are each what they would be as
separate processes, talking over sockets, and nothing below changes but the
URI. `just demo otel` runs them as separate processes.

1. **A broker** serves two streams, `logs` and `spans`, each with a log,
accepting remote publishers (`publish=True`).
1. **A broker** serves three streams, `logs`, `spans` and `metrics`, each
with a log, accepting remote publishers (`publish=True`).
2. **The producer**: two services, checkout and payments, handle traced
requests and log through the standard `logging` module. OTel's SDK turns
each into records and spans, published by the exporters in `logs.py` and
`spans.py`.
requests, log through the standard `logging` module, and count orders and
time requests. OTel's SDK turns each into records, spans and metric data
points, published by the exporters in `logs.py`, `spans.py` and
`metrics.py`.
3. **A live tail** follows the problems as they happen — `where=` a membership
filter on `severity_text` — and **a replay** reads back one request by its
trace id from both streams, `where={"trace_id": "<32 hex chars>"}`.
4. **The stored tables** answer what a log search and a trace view would:
errors per service, the slowest request, and ingest lag.
4. **The stored tables** answer what a log search, a trace view and a metrics
dashboard would: errors per service, the slowest request, ingest lag,
orders by outcome, mean request duration per service, and the failed
order's exemplar, which names the failed request's trace.
"""

from __future__ import annotations
Expand All @@ -33,7 +36,7 @@
from typing import TYPE_CHECKING, Any

import streamcast
from examples.otel import logs, spans
from examples.otel import logs, metrics, spans
from examples.otel.services import exporting, traffic

if TYPE_CHECKING:
Expand All @@ -42,10 +45,13 @@
# -- the broker -------------------------------------------------------------------


def streams(root: Path) -> tuple[streamcast.Stream, streamcast.Stream]:
def streams(
root: Path,
) -> tuple[streamcast.Stream, streamcast.Stream, streamcast.Stream]:
return (
streamcast.Stream.new("logs", root=root, schema=logs.SCHEMA),
streamcast.Stream.new("spans", root=root, schema=spans.SCHEMA),
streamcast.Stream.new("metrics", root=root, schema=metrics.SCHEMA),
)


Expand Down Expand Up @@ -75,9 +81,17 @@ async def one_trace(

async def main(root: Path) -> dict[str, Any]:
"""Run the whole demo under `root`, print it, and return what it saw."""
log_stream, span_stream = streams(root)
log_stream, span_stream, metric_stream = streams(root)
# No maintainer, as in the migration demo: a run this short never fills a
# log enough to seal it. Its five processes were still opening the logs
# when the run ended, and cost the run several seconds. `just demo otel`'s
# broker keeps `serve`'s default, which starts them.
server = await streamcast.serve(
[log_stream, span_stream], "127.0.0.1", 0, publish=True
[log_stream, span_stream, metric_stream],
"127.0.0.1",
0,
publish=True,
maintain=False,
)
base = f"ws://127.0.0.1:{server.sockets[0].getsockname()[1]}"
try:
Expand Down Expand Up @@ -116,13 +130,40 @@ async def main(root: Path) -> dict[str, Any]:
log_stream,
"SELECT max(streamcast_ts * 1000 - time_unix_nano) / 1e6 AS ms FROM log",
)
# Each row of `shop.orders` is one interval's count (delta temporality),
# so the total is a sum.
outcomes = query(
metric_stream,
"SELECT attributes['outcome'].string_value AS outcome, "
"sum(value_int)::BIGINT AS orders FROM log "
"WHERE name = 'shop.orders' GROUP BY outcome ORDER BY outcome",
)
durations = query(
metric_stream,
"SELECT service, sum(count)::BIGINT AS requests, "
"sum(sum) / sum(count) * 1e3 AS mean_ms FROM log "
"WHERE name = 'http.server.request.duration' "
"GROUP BY service ORDER BY service",
)
# The failed order was counted inside its request's span, so the SDK
# kept it as an exemplar: from the metric straight to the trace.
exemplars = query(
metric_stream,
"SELECT lower(hex(e.trace_id)) AS trace FROM ("
"SELECT unnest(exemplars) AS e FROM log WHERE name = 'shop.orders' "
"AND attributes['outcome'].string_value = 'failed')",
)
[stored_logs] = query(log_stream, "SELECT count(*) AS n FROM log")
[stored_spans] = query(span_stream, "SELECT count(*) AS n FROM log")
[stored_points] = query(metric_stream, "SELECT count(*) AS n FROM log")
finally:
server.close()
await server.wait_closed()

print(f"{stored_logs['n']} log records and {stored_spans['n']} spans stored")
print(
f"{stored_logs['n']} log records, {stored_spans['n']} spans and "
f"{stored_points['n']} metric data points stored"
)
print("live tail, errors and warnings as they happened:")
for message in live:
print(
Expand All @@ -142,6 +183,14 @@ async def main(root: Path) -> dict[str, Any]:
print(f"errors per service: {errors}")
print(f"slowest request: {slowest['ms']:.1f} ms")
print(f"ingest lag, worst: {lag['ms']:.1f} ms")
print(f"orders by outcome: {outcomes}")
for row in durations:
print(
f" {row['service']:>8} {row['requests']} requests, "
f"mean {row['mean_ms']:.2f} ms"
)

print(f"the failed order's exemplar: trace {[e['trace'] for e in exemplars]}")

return {
"logs": stored_logs["n"],
Expand All @@ -153,6 +202,10 @@ async def main(root: Path) -> dict[str, Any]:
"errors": errors,
"slowest_ms": slowest["ms"],
"lag_ms": lag["ms"],
"points": stored_points["n"],
"outcomes": outcomes,
"durations": durations,
"exemplars": [e["trace"] for e in exemplars],
}


Expand Down
Loading
Loading