A distributed data ingestion pipeline that fetches real-time weather telemetry for 500 geographical locations globally and stores it in a time-series database. Includes a real-time monitoring dashboard with REST APIs. Built with TypeScript/Node.js, Redis, and InfluxDB.
Note on location count: The pipeline is architected for 1,214 locations (240 named cities + 966 grid points at 8° resolution covering the entire globe). It is currently capped at 500 due to Open-Meteo's free tier quota of 10,000 req/day — at 500 locations/cycle × 1 cycle/min, the daily limit would be exceeded within minutes at full scale. To remove the cap, change
slice(0, 500)toslice(0, 1214)inservices/fetcher/src/locations.ts. The rate limiter, worker pool, and Redis stream architecture handle the full 1,214 locations without any other changes.
The dashboard provides real-time visibility into the entire pipeline: cycle progress, queue/stream depths, rate limiter state, E2E latency, and collected weather data — all in one view.
Open-Meteo API
│ HTTP (8 req/s, token bucket rate limiter — atomic Lua script in Redis)
▾
[Fetcher Service] (/metrics, /healthz)
node-cron scheduler — enqueues 500 locations every 60s
backpressure check — skips cycle if processor PEL count > threshold
50 async workers — BRPOP from queue, fetch, XADD to stream
│ XADD MAXLEN ~10,000
▾
Redis Stream (weather:raw) ←──────────────────────── Redis LIST (weather:locations:queue)
XREADGROUP (consumer group) scheduler LPUSH
Pending Entries List (PEL) for at-least-once delivery workers BRPOP
│
▾
[Processor Service]
reads messages → writes to InfluxDB → XACK
batched writes: 100 points/s or 1s flush interval
│
▾
InfluxDB (weather_bucket, 30-day retention)
measurement: weather | tags: city_name, weather_condition
fields: temperature, latitude, longitude | timestamp: recorded_at
[Dashboard Service] (port 4000)
reads Redis → pipeline status, stream depth, rate limiter state
reads InfluxDB → latest readings, city history, summary stats
scrapes /metrics → Prometheus counters (success rate, latency)
serves Web UI → http://localhost:4000
| Decision | Why |
|---|---|
| 50 async workers (not threads) | Weather fetch is I/O-bound — async coroutines yield while waiting on TCP, zero CPU overhead |
| Token bucket in Redis (Lua) | 50 workers share one rate limit; Lua script makes consume+refill atomic, no TOCTOU race |
| Redis Streams (not List/PubSub) | Consumer groups give PEL-based at-least-once delivery; messages survive processor crashes |
| PEL count for backpressure | XLEN counts all messages including already-processed ones; PEL counts only unfinished work |
| InfluxDB over PostgreSQL | Native time-series model, built-in retention, Flux aggregations — purpose-built for this workload |
| Separate dashboard service | Needs read access to both Redis and InfluxDB; neither fetcher nor processor has both |
MAXLEN ~10000 (approximate) |
Exact trim is O(N) per write; ~ amortizes cost across bulk node trims — bounded memory |
At 500 locations/cycle, 8 req/s rate limit:
| Metric | Value |
|---|---|
| Cycle duration | ~62 seconds (500 ÷ 8 req/s) |
| Throughput | ~8 weather readings/second |
| Stream peak depth | ~500–10,000 entries (processor catches up within seconds) |
| InfluxDB write rate | up to 100 points/batch, 1s flush |
| E2E latency (steady state) | 300–800 ms (queue wait + stream wait + InfluxDB write) |
| Redis memory | bounded — stream capped at ~10,000 entries ≈ 2MB |
Responsible for fetching weather data from the Open-Meteo API and publishing it to a Redis Stream.
Scheduler (scheduler.ts)
- Uses
node-cronto enqueue all 500 locations into a Redis LIST every 60 seconds - Backpressure: before each cycle, checks the processor group's PEL (Pending Entries List) count via
XINFO GROUPS. If un-ACKed messages exceed the threshold (default 5000), the cycle is skipped. PEL count is the correct signal — unlikeXLEN, it only counts messages the processor hasn't finished, not historical processed entries. - First cycle always runs unconditionally — ensures the pipeline starts immediately
- Each cycle gets a monotonic ID and a start timestamp stored in Redis, used by the analytics reporter
Workers (worker.ts)
- 50 concurrent async workers, each running an infinite loop:
BRPOP→acquire token→fetch→XADD - Stream trimming:
XADDusesMAXLEN ~ 10000to cap the stream at approximately 10,000 entries. The~flag lets Redis trim in efficient batches rather than exact per-write trimming. Prevents unbounded memory growth. - All 50 start simultaneously — the rate limiter handles concurrency, not startup staggering
- Per-second analytics reporter prints live cycle progress: requests/sec, success/fail counts, avg + p99 latency
Rate Limiter (rate-limiter.ts)
- Token bucket algorithm implemented via an atomic Redis Lua script
- Lua ensures no two workers can "double-spend" the same token even with 50 concurrent async workers hitting Redis simultaneously
- Cap: 8 tokens/sec — comfortably under Open-Meteo's 600 req/min free tier limit
- Cooldown mode: if a 429 slips through, a Redis key with a TTL is set — all workers sleep until it expires (
PTTLfor exact sleep duration, not polling) - Backed by Redis so it works correctly across multiple fetcher replicas in production
Fetcher (fetcher.ts)
- Shared axios instance with IPv4 forced (
family: 4) — avoids IPv6 DNS hangs on some network configs axiosRetrywith full-jitter exponential backoff, up to 5 retries on 5xx and network errors- Respects
Retry-Afterheader from server responses - 429 (rate limit) is not retried — handled at the worker level by
notifyThrottled()which triggers a 30s global cooldown USE_MOCK=trueswaps the real HTTP call formockFetchWeather()— same return type, no quota used
Mock Weather (mock-weather.ts)
- Drop-in replacement for the real API — produces realistic temperatures based on latitude and season
- Simulates network latency (80–350ms) so worker timing matches real conditions
- Used for local development when the Open-Meteo daily quota is exhausted
Metrics & Server (metrics.ts, server.ts)
- Prometheus counters:
api_calls_total,api_calls_success_total,api_calls_failed_total,rate_limiter_denials_total - Prometheus histogram:
api_response_latency_seconds(buckets: 50ms → 10s) - Express server on port 3000:
GET /metrics— Prometheus scrape endpointGET /healthz— probes Open-Meteo reachability, returns200 okor503 degraded
A read-only microservice that provides REST APIs and a web dashboard for monitoring the pipeline and exploring collected weather data. Separated from the fetcher and processor because it needs read access to both Redis (pipeline state) and InfluxDB (weather data queries) — neither the fetcher nor the processor has both.
REST API Endpoints:
| Endpoint | Description |
|---|---|
GET /api/pipeline/status |
Cycle ID, queue depth, stream length, pending count, rate-limiter cooldown, backpressure state |
GET /api/weather/latest |
Latest temperature + condition for each city (optional ?city= filter) |
GET /api/weather/summary |
Global stats: total cities, avg/min/max temp, most common condition |
GET /api/weather/:city/history?range=6h |
Temperature time-series for one city (ranges: 1h, 6h, 12h, 1d, 3d, 7d) |
GET /api/stream/recent?count=20 |
Last N raw entries from the Redis stream |
Web Dashboard (public/):
- Vanilla HTML/CSS/JS — no build step, no frontend framework
- Dark theme, responsive layout, polls APIs every 5 seconds
- Pipeline status bar with color-coded backpressure indicator (green/yellow/red)
- Summary cards (cities count, avg/min/max temp, top condition)
- Sortable, searchable weather data table — click any city to drill into its history
- City detail panel with a Chart.js temperature sparkline (selectable time range)
- Live stream feed showing raw data flowing through the pipeline
Reads from the Redis Stream and writes to InfluxDB.
Consumer (consumer.ts)
- Uses Redis
XREADGROUPwith a named consumer group (processor-group) - Consumer name defaults to
processor-{hostname}— unique per instance, enabling horizontal scaling via multiple processor replicas sharing the same group - On startup, first reads with
"0"to drain any pending (unacknowledged) messages from before a crash — ensures no data loss on restart - Then switches to
">"for new messages XACKis sent after the InfluxDB write is enqueued (not after flush). Crash within the 1s flush window may lose up to 1s of data — acceptable for telemetry. Reprocessed messages on restart are safe due to InfluxDB idempotency.- Exponential backoff (up to 30s) kicks in after 5 consecutive write failures — prevents tight-looping during InfluxDB outages
InfluxDB Writer (influx-writer.ts)
- Connects to InfluxDB using
@influxdata/influxdb-client - Batches writes: flushes every 1 second or when 100 points accumulate — more efficient than one write per point
- Each point uses
recorded_atas the timestamp — this is the actual observation time from Open-Meteo, not the ingestion time - Idempotency: InfluxDB deduplicates by
(measurement + tags + timestamp). Writing the same point twice (e.g. after a processor restart) overwrites rather than duplicates — no UPSERT logic needed
Acts as three things simultaneously:
- Job queue — a LIST (
weather:locations:queue) that the scheduler pushes to and workers pop from - Stream broker — a Stream (
weather:raw) that workers publish to and the processor consumes from - Rate limiter state — stores token bucket state (
rate_limiter:weather_api:bucket) and cooldown flag (rate_limiter:weather_api:cooldown)
Chosen over PostgreSQL + TimescaleDB for this use case because:
- Native time-series data model — no schema migrations needed when adding fields
- Built-in retention policies — 30-day automatic data expiry configured at bucket creation
- Flux query language is purpose-built for time-series aggregations (windowed averages, downsampling)
- Free tier sufficient for development and demonstration
Data model:
measurement: weather
tags: city_name (string), weather_condition (string)
fields: temperature (float), latitude (float), longitude (float)
timestamp: recorded_at (ms precision)
Tags are indexed — used for filtering and grouping. Fields are the measured values.
Prerequisites: Docker, Docker Compose
# Start all five services (Redis + InfluxDB + fetcher + processor + dashboard)
docker compose up --build
# Use mock data instead of real API (no quota consumed)
# Edit docker-compose.yml, uncomment: USE_MOCK: "true"
# Then restart: docker compose up --build
# Stop
docker compose down
# Stop and wipe all persisted data (Redis + InfluxDB volumes)
docker compose down -vEndpoints once running:
| URL | Description |
|---|---|
| http://localhost:4000 | Dashboard UI |
| http://localhost:4000/api/pipeline/status | Pipeline status API |
| http://localhost:4000/api/weather/latest | Latest weather data API |
| http://localhost:8086 | InfluxDB UI (admin / adminpassword) |
| http://localhost:3001/metrics | Prometheus metrics |
| http://localhost:3001/healthz | Health check |
Open InfluxDB UI → Data Explorer → Script Editor.
Count of cities written in the last hour:
from(bucket: "weather_bucket")
|> range(start: -1h)
|> filter(fn: (r) => r._measurement == "weather" and r._field == "temperature")
|> group()
|> count()
Temperature time series for a specific city:
from(bucket: "weather_bucket")
|> range(start: -24h)
|> filter(fn: (r) => r._measurement == "weather" and r._field == "temperature")
|> filter(fn: (r) => r.city_name == "London")
|> sort(columns: ["_time"])
Full table — one row per observation, all fields:
from(bucket: "weather_bucket")
|> range(start: -1h)
|> filter(fn: (r) => r._measurement == "weather")
|> pivot(rowKey: ["_time", "city_name"], columnKey: ["_field"], valueColumn: "_value")
|> keep(columns: ["_time", "city_name", "weather_condition", "temperature", "latitude", "longitude"])
Observations per minute (confirms cycle timing):
from(bucket: "weather_bucket")
|> range(start: -2h)
|> filter(fn: (r) => r._measurement == "weather" and r._field == "temperature")
|> aggregateWindow(every: 1m, fn: count, createEmpty: false)
Prerequisites: kubectl, a running cluster (minikube, GKE, EKS, etc.)
# Deploy everything
kubectl apply -f k8s/
# Check pod status
kubectl get pods
# Stream logs
kubectl logs -f deployment/fetcher
kubectl logs -f deployment/processor
# Tear down
kubectl delete -f k8s/K8s manifest overview:
| File | What it creates |
|---|---|
configmap.yaml |
Shared env vars — Redis URL, InfluxDB URL, org, bucket |
secret.yaml |
Sensitive values — InfluxDB token + admin password |
redis-deployment.yaml |
Redis Deployment + internal Service (redis-service:6379) |
influxdb-deployment.yaml |
InfluxDB Deployment + internal Service (influxdb-service:8086) |
fetcher-deployment.yaml |
Fetcher Deployment + Service (exposes port 3000 for metrics) |
processor-deployment.yaml |
Processor Deployment |
dashboard-deployment.yaml |
Dashboard Deployment + Service (port 4000) |
Services use K8s internal DNS — pods reach each other by service name, not IP.
Open-Meteo free tier limits: 600 req/min, 10,000 req/day.
| Value | |
|---|---|
| Pipeline rate | 8 req/s = 480 req/min |
| Per cycle | 500 locations |
| Cycle duration | ~62 seconds at 8 req/s |
| Daily usage (1 cycle/min) | ~720,000 req — exceeds free tier |
| Safe usage | Run in mock mode for development; use real API for demos only |
For sustained production use, either purchase the Open-Meteo commercial plan ($29/month for 10M req/month) or reduce cycle frequency to once per hour (500 req/hr = 12,000 req/day).
| Variable | Service | Default | Description |
|---|---|---|---|
REDIS_URL |
fetcher, processor, dashboard | redis://localhost:6379 |
Redis connection URL |
INFLUX_URL |
processor, dashboard | http://localhost:8086 |
InfluxDB connection URL |
INFLUX_TOKEN |
processor, dashboard | my-super-secret-token |
InfluxDB API token |
INFLUX_ORG |
processor, dashboard | weather_org |
InfluxDB organisation |
INFLUX_BUCKET |
processor, dashboard | weather_bucket |
InfluxDB bucket name |
USE_MOCK |
fetcher | unset | Set to "true" to use mock data |
METRICS_PORT |
fetcher | 3000 |
Port for /metrics and /healthz |
PORT |
dashboard | 4000 |
Dashboard server port |
BACKPRESSURE_THRESHOLD |
fetcher | 5000 |
PEL count at which fetcher skips a cycle |
CONSUMER_NAME |
processor | processor-{hostname} |
Redis stream consumer identity |
STREAM_MAXLEN |
dashboard | 10000 |
Expected stream max length shown in UI |