Modular, production-minded real-time telemetry processing pipeline in Python. Ingests high-frequency sensor data, validates and enriches events, computes windowed aggregations, stores time-series data, detects anomalies with an online ML ensemble, and visualizes results in Grafana or Streamlit.
flowchart LR
subgraph Sources
SIM[Sensor Simulator]
MQTT[MQTT Broker]
KAFKA[Kafka / Redpanda]
WS[WebSocket Clients]
end
subgraph Pipeline
ING[Ingestion Layer]
VAL[Schema Validation + Dedup]
ENR[Enrichment]
WIN[Windowed Aggregation]
ANO[Anomaly Ensemble]
ALT[Alert Dispatcher]
end
subgraph Storage
TS[(TimescaleDB)]
MEM[(Memory - tests)]
end
subgraph Viz
API[HTTP API :8080]
ST[Streamlit Dashboard]
GF[Grafana]
end
SIM --> WS
MQTT --> ING
KAFKA --> ING
WS --> ING
ING --> VAL --> ENR --> WIN
ENR --> ANO --> ALT
ENR --> TS
WIN --> TS
ANO --> TS
TS --> API --> ST
TS --> GF
- Multi-transport ingestion: MQTT, Kafka, WebSocket (config switch)
- Validation: JSON Schema + per-sensor range checks + deduplication
- Stream processing: Tumbling/sliding windows with mean/min/max/std/count
- Storage: TimescaleDB (production) or in-memory (tests)
- Anomaly detection ensemble:
- Statistical (EWMA + z-score)
- Online HalfSpaceTrees (
river) - Online autoencoder (numpy / PyTorch / ONNX)
- Rule-based thresholds from config
- ADWIN concept drift detection
- Live API: HTTP API on
:8080with JSON metrics for Streamlit - Prometheus:
/metricsendpoint + Prometheus server in Docker Compose - Benchmark harness: Formal throughput/latency report (
telemetry-benchmark) - Synthetic data: High-frequency generator with labeled anomaly injection
- Dataset replay: NAB-style and pump CSV samples
- Testing: Unit, integration, load (1k events), and latency benchmarks
- Python 3.11+
- Docker & Docker Compose (for full stack)
cd telemetry-pipeline
python -m venv .venv && source .venv/bin/activate
pip install -e ".[dev]"
# Terminal 1 — pipeline (WebSocket ingestion, in-memory storage for quick demo)
# Edit config/pipeline.yaml: set storage.backend: memory
python -m telemetry.main --config config/pipeline.yaml
# Terminal 2 — simulator
python -m telemetry.simulator.generator --duration 60
# Terminal 3 — product dashboard (Next.js, polls /api/metrics)
cd frontend && npm install && npm run dev
# Legacy: streamlit run src/telemetry/viz/streamlit_app.py
# Benchmark harness
telemetry-benchmark --events 10000 --report benchmark_report.jsondocker compose up --build| Service | URL / Port |
|---|---|
| Pipeline API | http://localhost:8081/api/metrics |
| Prometheus | http://localhost:9090 |
| Frontend | http://localhost:3001 (Signal dashboard) |
| Grafana | http://localhost:3000 (admin/admin) |
| TimescaleDB | localhost:5432 |
| Kafka | localhost:19092 |
| MQTT | localhost:1883 |
The simulator starts automatically and feeds the pipeline for 1 hour.
Login at http://localhost:3000 (admin/admin) → Dashboards → Telemetry → Telemetry Pipeline Overview
The dashboard includes 14 panels:
| Section | Panels |
|---|---|
| Overview stats | Events/hr, active devices, anomalies/hr, avg ingest latency |
| Ingestion | Event rate, events by sensor type |
| Sensor metrics | Industrial temperature, windowed vibration |
| Latency | Ingest latency avg/P95 (Timescale), processing latency (Prometheus) |
| Anomalies | Detection rate, severity breakdown, recent anomalies table |
| Pipeline health | Throughput (eps) from Prometheus |
Datasources auto-provisioned: TimescaleDB (events/anomalies) + Prometheus (pipeline self-metrics).
To reload after editing docker/grafana/provisioning/dashboards/telemetry-overview.json:
docker compose restart grafanaAll behavior is driven by YAML:
| File | Purpose |
|---|---|
config/pipeline.yaml |
Transport, validation, processing, storage, anomaly, alerting |
config/sensors.yaml |
Sensor type definitions, baselines, rules |
config/schemas/sensor_event.json |
JSON Schema for events |
- Create a Slack Incoming Webhook
- Copy
.env.example→.envand set:
TELEMETRY_ALERTING_ENABLED=true
TELEMETRY_SLACK_WEBHOOK_URL=https://hooks.slack.com/services/...
TELEMETRY_ALERTING_MIN_SEVERITY=medium
TELEMETRY_ALERTING_COOLDOWN_SECONDS=60- Restart the pipeline:
docker compose up -d pipeline
Alerts include severity color, device, score, top detection method, and drift flag. Cooldown prevents duplicate alerts per device.
ingestion:
transport: kafka # mqtt | kafka | websocket- Add entry under
sensor_typesinconfig/sensors.yaml - Add enum value in
config/schemas/sensor_event.json - Optionally add rule thresholds under
rules
Example event:
{
"device_id": "industrial-device-001",
"sensor_type": "industrial",
"timestamp": "2026-06-17T12:00:00Z",
"metrics": {"temperature": 65.2, "pressure": 4.5, "vibration": 3.1},
"tags": {"line": "A3"}
}# Run pipeline
telemetry-pipeline --config config/pipeline.yaml
# Synthetic sensor data generator
telemetry-simulator --devices 10 --interval-ms 50 --anomaly-rate 0.05
# Replay labeled CSV dataset
telemetry-replay --csv data/sample/pump_sample.csv --speed 50
# Streamlit dashboard
telemetry-dashboard
# Throughput/latency benchmark
telemetry-benchmark --events 10000 --warmup 500
# Load test (100k+ eps target)
telemetry-load --mode direct --events 100000 --target-eps 100000
MODE=kafka-producer EVENTS=1000000 ./scripts/run-load-test.shpip install -e ".[dev]"
pytest -v
pytest tests/test_load.py -v # 1k-event throughput smoke
pytest tests/test_load_100k.py -v # load harness + 5k direct soak
pytest tests/test_latency.py -v # latency budget check
RUN_SLOW_LOAD=1 pytest -m slow -v # optional 100k-event soakThe pipeline validates configuration at startup (before connecting to brokers/DB):
- Ensemble weights sum to ~1.0
- JSON schema file exists and matches
sensors.yamltypes - Window/slide interval consistency
- Alerting enabled → webhook URL required
- Optional
--strictflag fails if TimescaleDB is unreachable
telemetry-pipeline --config config/pipeline.yaml --strictOn every push/PR to main, CI runs:
- test — full
pytestsuite - smoke —
docker compose up, verifies/healthand events flowing on:8081
Run the smoke test locally:
chmod +x scripts/smoke_test.sh
./scripts/smoke_test.shtelemetry-pipeline/
├── config/ # YAML schemas, sensor definitions
├── data/sample/ # NAB + pump CSV samples
├── docker/ # Dockerfiles, Timescale init, Grafana provisioning
├── src/telemetry/
│ ├── ingestion/ # MQTT, Kafka, WebSocket
│ ├── validation/ # Schema + enrichment
│ ├── processor/ # Windowed aggregation
│ ├── storage/ # TimescaleDB + memory
│ ├── anomaly/ # Ensemble detector + drift
│ ├── simulator/ # Generator + replay
│ └── viz/ # Streamlit + HTTP API
└── tests/ # Unit, integration, load, latency
- Create
src/telemetry/anomaly/my_method.py - Wire into
AnomalyDetector.detect()indetector.py - Add weight in
config/pipeline.yamlunderanomaly.ensemble_weights
anomaly:
autoencoder:
backend: torchInstall ML extras: pip install -e ".[ml]"
from telemetry.anomaly.autoencoder import export_numpy_to_onnx
export_numpy_to_onnx(input_dim=3, hidden_dim=8, output_path="models/autoencoder.onnx")Then set anomaly.autoencoder.backend: onnx in config.
- Implement
StorageBackendinsrc/telemetry/storage/clickhouse.py - Register in
create_storage()factory - Add
storage.backend: clickhouseoption
- Run pipeline + simulator as separate K8s deployments
- Use Kafka for ingestion at scale
- Enable alerting with Slack webhook in config
- Add Prometheus metrics exporter (hook into
PipelineMetrics) - Train custom autoencoder and serve via ONNX Runtime
| Metric | Local (memory) | Docker (Timescale) | Load test target |
|---|---|---|---|
| Throughput | >1,000 eps | 500+ eps | 100,000+ eps (Kafka producers / scaled K8s) |
| Per-event processing | <100ms | <100ms | sampled in telemetry-load |
| End-to-end ingest | <50ms typical | <100ms | e2e-kafka mode |
Use config/pipeline.load.yaml (anomaly/Prometheus off, count-only memory storage):
# Direct: in-process minimal ingest path (multi-core via --workers)
telemetry-load --mode direct --events 100000 --workers 8 --target-eps 100000
# Kafka producer flood — primary path to 100k+ eps (needs Redpanda/Kafka)
telemetry-load --mode kafka-producer --events 1000000 --workers 8 --target-eps 100000
# End-to-end: producers + pipeline consumer (full validation path)
telemetry-load --mode e2e-kafka --duration 30 --workers 8 --target-eps 100000
./scripts/run-load-test.sh # MODE=direct|kafka-producer|e2e-kafka| Mode | What it measures |
|---|---|
direct |
Single-host ingest ceiling (process_event_minimal, count-only storage) |
kafka-producer |
Multi-process publish rate to Kafka (100k+ eps with enough workers/partitions) |
e2e-kafka |
Producer + consumer throughput with full pipeline features |
Kubernetes: kubectl apply -f k8s/load-test-job.yaml (after pipeline + Redpanda are up).
Scale pipeline HPA replicas to consume sustained Kafka load.
Reports are written to load_test_report.json with producer/consumer eps and latency percentiles.
MIT