StreamAlpha
A real-time distributed market analytics platform built around event streaming, low-latency state, and WebSocket delivery capable of 10,000 events/second.
Real-Time Event Stream & Backpressure Pipeline
Producer → Kafka/Redpanda → Stream Processor → Redis/PostgreSQL → WebSocket Clients
Tick Producer
Financial market feeds (AAPL, NVDA, MSFT) publishing live tick events.
Kafka / Redpanda
Partitioned topic logs guaranteeing per-symbol total order under load.
FastAPI Gateway
Asynchronous ring buffer manager dispatching frames & handling backpressure.
Redis + Postgres
Redis for sub-1ms state snapshot cache; PostgreSQL for 1s/1m OHLCV batch aggregations.
WASM Clients
Perspective WebAssembly rendering in Web Worker at 60 FPS without freezing UI.
01. Problem Statement & Motivation
High-frequency financial market updates overwhelm typical web architectures: WebSocket broadcast storms freeze browser client renderers, slow consumers cause catastrophic memory leaks in the backend, and database write throughput collapses when persisting sub-second tick streams.
02. System Architecture Design
Engineered an event-driven decoupled pipeline. Market tick producers broadcast to Kafka/Redpanda partitions. An independent FastAPI WebSocket gateway subscribes to partitioned topics, manages per-client bounded ring buffers, drops slow-consumer frames under configurable backpressure, and caches the latest market snapshot in Redis. An independent analytics consumer aggregates 1-second and 1-minute OHLCV candles for batch persistence into PostgreSQL.
- Slow-Client Message Eviction: Drops non-critical market ticks when client buffer reaches capacity, sending an explicit lag warning packet.
- Symbol-Based Topic Partitioning: Ensures tick ordering per equity symbol across Kafka consumer groups.
- Prometheus Telemetry Instrumentation: Exposes live metrics for queue buffer depth, drop rates, Kafka consumer lag, and end-to-end latency.
- Automated Load Testing Harness: Configurable benchmark suite verifying throughput from 100 to 10,000 events/second.
03. Architectural Decisions & Tradeoffs
Per-Client Bounded Asynchronous Ring Queues
Each WebSocket client connection receives an isolated bounded queue. If a client stalls, oldest non-critical ticks are evicted without blocking the shared event ingestion pipeline.
FINOS Perspective WebAssembly Engine in Web Worker
Processed 10,000 updates/second entirely in a Web Worker running C++ compiled to WebAssembly. Only viewport diffs are sent to the DOM, keeping main thread UI rendering rock solid at 60 FPS.
Independent Analytics Consumer with Micro-Batching
Decoupled live WebSocket delivery from relational persistence. Tick events are aggregated in-memory over 1-second sliding windows before executing PostgreSQL COPY batch inserts.
04. Verified Empirical Outcomes
| Metric Dimension | Guarded Platform | Significance |
|---|---|---|
| 100 evt/s Throughput | p50: 2.70 ms | p95: 7.10 ms | Baseline low-frequency streaming with pristine delivery |
| 1,000 evt/s Throughput | p50: 2.57 ms | p95: 6.87 ms | Standard market hours trading load |
| 5,000 evt/s Throughput | p50: 5.83 ms | p95: 10.79 ms | High volatility surge simulation |
| 10,000 evt/s Burst Load | p50: 12.56 ms | p95: 17.09 ms | Sustained peak market opening stress test at 60 FPS UI |
05. Production Roadmap & Next Iterations
- >Implement zero-copy serialization using Apache Arrow Flight for cross-service analytics streaming.
- >Add kernel-bypass DPDK networking for sub-microsecond tick capture on bare metal.