Low-Latency WebSocket Aggregator for Market Data
Trading bots, market makers, and risk systems depend on fresh market data—a delay of milliseconds can cost profits. In one project, a client lost tens of thousands of dollars in a month due to a stale connection: data stopped flowing, but the system continued trading based on outdated prices. Each exchange uses its own protocol, connection limits, and message format. A low-latency crypto data aggregator solves this.
According to the WebSocket API documentation, Binance allows up to 1024 streams per connection, while Bybit only 10 topics. Building a universal aggregator that connects to all exchanges, normalizes streams, and delivers data with minimal latency is a non-trivial task. Even an error in stream processing can lead to arbitrage losses or incorrect order execution. We built a modular aggregator that addresses these issues. Our solutions have been proven in projects handling 100,000 messages per second on a single core using Python asyncio—an order of magnitude faster than standard multithreaded approaches.
We have delivered over 50 projects in crypto infrastructure. Below we break down the key components and architecture.
How the Aggregator Handles Different Exchange Limits
The Connection Manager automatically distributes subscriptions, respecting each exchange’s limits. For each exchange we configure a manager that creates new connections when the limit is exhausted.
| Exchange |
Max streams / conn |
Ping interval |
Max connections |
| Binance |
1024 |
3 min |
Unlimited |
| Bybit |
10 topics / conn |
20 sec |
Unlimited |
| OKX |
240 channels / conn |
30 sec |
Unlimited |
| Kraken |
Not documented |
Adaptive |
Unlimited |
class ConnectionManager:
def __init__(self, max_per_conn: int = 900):
self.connections: list[WSConnection] = []
self.max_per_conn = max_per_conn
self.subscriptions: dict[str, WSConnection] = {}
async def subscribe(self, channels: list[str]):
for channel in channels:
conn = self._find_or_create_connection()
await conn.subscribe(channel)
self.subscriptions[channel] = conn
def _find_or_create_connection(self) -> WSConnection:
for conn in self.connections:
if conn.subscription_count < self.max_per_conn:
return conn
new_conn = WSConnection(self.on_message, self.on_disconnect)
self.connections.append(new_conn)
return new_conn
async def on_disconnect(self, conn: WSConnection):
# Exponential backoff and resubscribe
await asyncio.sleep(conn.backoff.next())
await conn.reconnect()
await conn.resubscribe()
When the limit is exceeded, the Manager automatically creates an additional connection. For example, for Bybit with a limit of 10 topics per connection, subscribing to 25 channels will result in 3 connections. Exponential backoff prevents exchange overload during mass disconnects.
Why Heartbeat and Stale Detection Matter
Exchanges may go silent without a TCP disconnect—the connection is alive but no data arrives. A watchdog timer for each connection solves this. If no message arrives for more than 30 seconds, the connection is forcibly recreated. Heartbeat monitoring and stale detection are key elements of a robust aggregator.
class HeartbeatMonitor:
STALE_THRESHOLD_SEC = 30
async def watch(self, conn: WSConnection):
while True:
await asyncio.sleep(5)
age = time.time() - conn.last_message_time
if age > self.STALE_THRESHOLD_SEC:
logger.warning(f"Stale connection detected, forcing reconnect")
await conn.force_reconnect()
In the project I mentioned, the lack of such a monitor led to losses. After deploying the aggregator with the Heartbeat monitor, incidents stopped, and the savings on missed profit reached about 40%.
Publishing Data to Consumers
The aggregator publishes normalized data through multiple channels. The choice depends on reliability and latency requirements.
| Channel |
Latency |
Reliability |
Persistence |
Typical use‑case |
| Redis Pub/Sub |
<1 ms |
No guarantee |
No |
Real‑time broadcast without log |
| Redis Streams |
<5 ms |
Guaranteed (consumer groups) |
Yes |
Recovery after downtime |
| Kafka streaming |
<10 ms |
Guaranteed (commit log) |
Yes |
High‑load systems |
| gRPC streaming |
<1 ms |
Guaranteed (bidirectional) |
No |
Direct client‑aggregator connection |
Redis Pub/Sub offers minimal latency but no delivery guarantee. Redis Streams and Kafka are suitable for reliable delivery with ability to replay missed messages. gRPC streaming is for direct low‑latency connections.
Performance Metrics
The aggregator exports Prometheus metrics:
-
ws_messages_received_total{exchange, channel}
-
ws_message_latency_ms{exchange}
-
ws_reconnects_total{exchange}
-
ws_active_connections{exchange}
-
ws_subscription_count{exchange}
These metrics allow quick identification of connection issues and overloads. With a proper implementation in Python (asyncio), the aggregator processes 50,000–100,000 messages per second on a single core. Go or Rust can handle an order of magnitude more.
Process and Workflow
-
Analysis – we study the list of exchanges, data types (order book, trades, tickers), and latency requirements.
-
Design – we choose the stack (Python/Go, Redis/Kafka) and design the normalization schema.
-
Implementation – we build the Connection Manager, Heartbeat, and publication modules.
-
Testing – we simulate disconnects, run load tests, and verify recovery.
-
Deployment – we deploy in your infrastructure (k8s, bare metal) and configure monitoring.
Typical mistakes in DIY implementations include ignoring exchange limits—exceeding max streams causes disconnection; lack of heartbeat—stale connections lead to trading on outdated data; synchronous processing—blocking calls kill performance; absence of metrics—impossible to assess system health. Our aggregator solves each of these problems.
Timelines and Cost
A basic aggregator for one exchange takes 2–4 weeks. Adding an additional exchange takes 1–2 weeks. A full solution with Kafka and dashboards starts from 2 months. Cost is calculated individually. Infrastructure savings compared to purchased solutions can reach 40%. Contact us for an assessment of your project—we will prepare a proposal within 1–2 days. Get a consultation on your data pipeline architecture. Order development of a WebSocket aggregator for your trading system.
Why exchange development requires deep domain expertise
We develop exchanges — not 'chart sites,' but matching engines that process thousands of orders per second without delay, route liquidity between pools, and guarantee that no user gains access to others' funds. Teams that start with the UI and postpone the engine 'for later' end up rewriting everything in six months in 90% of cases.
Order Book vs AMM: where most projects break
Centralized exchanges (CEX) are built around an order book + matching engine. Decentralized exchanges (DEX) either also use an order book (dYdX on StarkEx, Serum/OpenBook on Solana) or an AMM with concentrated liquidity (Uniswap v3/v4, Curve, Balancer). A classic mistake when developing a CEX is implementing the matching engine on top of a relational database with transactions for each match. PostgreSQL handles ~500 RPS without special effort, but at peak loads of 5,000–10,000 orders per second, it turns into a deadlock nightmare. The correct architecture: in-memory order book (Redis Sorted Sets or custom C++/Rust structure), asynchronous writing of matches to PostgreSQL via a queue (Kafka/RabbitMQ), and a separate settlement service that finally updates balances.
For DEX, the most painful problem is sandwich attacks and MEV. A pool with a plain xy=k AMM without slippage protection becomes a target for MEV bots within hours of launch. Uniswap v2 lost hundreds of millions of dollars in user liquidity. Solutions: integration with Flashbots Protect, a commit-reveal scheme for orders, or switching to TWAMM (Time-Weighted AMM) for large trades.
Concentrated liquidity and impermanent loss
Uniswap v3 introduced concentrated liquidity – LPs choose a price range in which to provide liquidity. Capital efficiency increased 4,000x compared to v2 for stable pairs. But implementing this mechanism correctly is non-trivial. The Uniswap v3 liquidity contract uses tick-based accounting: the price space is divided into discrete ticks (tick = log₁.0001(price)), each tick stores accumulated fee growth and liquidity delta. When creating a position, the lower and upper ticks are computed, and the contract recalculates all active positions at each swap. Storage layout is critical here – incorrect variable packing in slots easily adds 40–60% to swap gas cost.
We implemented a Uniswap v3 fork for a client on Polygon with a custom fee tier system. The initial version consumed 180k gas for a swap across 2 ticks. After slot packing of variables in Tick.Info and inlining several internal calls, it dropped to 112k gas. This reduced gas costs by 38% and saved the client substantial costs on fees monthly. The techniques applied are described in the Uniswap v3 Whitepaper and confirmed by our audit experience.
How a matching engine delivers performance
A production-ready matching engine is built according to the following scheme:
-
Order ingestion layer – WebSocket gateway (Go or Rust), accepts orders, validates signature, checks balance via Redis, queues them. Latency at this level must be <1ms.
-
Matching core – single-threaded event loop (eliminates race conditions without mutexes). In memory, we hold two Sorted Sets for each trading instrument: bids and asks. FIFO matching for limit orders, immediate-or-cancel for market orders. Throughput with a proper Rust implementation – 500k–1M matches per second on a single core.
-
Settlement service – reads matches from Kafka, atomically updates balances in PostgreSQL (
UPDATE accounts SET balance = balance - $1 WHERE id = $2 AND balance >= $1). Optimistic locking via row versioning.
-
Withdrawal pipeline – separate service with cold/hot wallet architecture. The hot wallet holds 5–10% of total deposits, the rest is cold storage with multi-sig (Gnosis Safe or custom HSM). Automatic withdrawals only from hot wallet, large amounts require manual authorization.
| Component |
Technology |
Latency / Throughput |
| Order gateway |
Go + WebSocket |
<1ms p99 |
| Matching engine |
Rust (in-memory) |
500k+ orders/sec |
| Balance store |
Redis (write-through) |
<0.5ms |
| Settlement DB |
PostgreSQL 14+ |
~50k TPS with partitioning |
| Event streaming |
Apache Kafka |
1M+ events/sec |
| Blockchain node |
Geth / Solana validator |
depends on chain |
How our exchange development process ensures reliability
Smart contracts and gas optimization
For EVM-based DEX (Ethereum, Arbitrum, Optimism, Polygon), the entire critical path lives in Solidity. Main contracts: Pool, Factory, Router, PositionManager (for v3-like), and Quoter for off-chain calculations. Typical mistakes we see in audits:
Reentrancy via callback. Uniswap v3 uses flash swap with a callback (uniswapV3SwapCallback). If your router lacks a nonReentrant guard and you don't check msg.sender == pool, the contract gets drained via a nested call. This is not hypothetical – several v3 forks lost funds this way.
Oracle manipulation in AMM. If your contract uses the spot price from the pool for collateral calculation, it is front-runnable. Correct: TWAP over 30+ minutes (Uniswap v3 OracleLib) or an external oracle (Chainlink).
Unbounded loops in liquidity range. If a swap crosses many ticks in a row (price impact 80%+), gas may exceed the block limit. Need MAX_TICKS_CROSSED with partial fill and returning the remainder.
For Solana DEX (Anchor framework, Rust), the architecture is fundamentally different: account-based model, Program Derived Addresses (PDA) instead of storage, Cross-Program Invocations instead of internal calls. Solana's throughput (~3,000–4,000 TPS vs 15–30 on Ethereum mainnet) allows building on-chain order books – exactly what Phoenix DEX does.
Liquidity bootstrapping and aggregator integration
Launching a pool is not enough – you need to ensure liquidity at launch. Practical mechanisms:
-
Liquidity Bootstrapping Pool (LBP) – initial price is high, asset weights dynamically shift, creating selling pressure and even token distribution. Implemented in Balancer v2.
-
Initial Liquidity Offering via Uniswap v3 – adding liquidity in a narrow range around the initial price, then gradually expanding as volume grows. Requires active liquidity management or integration with Arrakis/Gamma.
-
Integration with 1inch, Paraswap, Li.Fi – aggregators bring traffic but require standard compliance: the pool must have correct
getAmountsOut, support ERC-20 approval/permit, and not have custom transfer hooks that break the aggregator's routing.
Development process and deliverables
Analytics and design begin with choosing the architectural model: CEX with custodial storage, non-custodial DEX, or hybrid (off-chain order book + on-chain settlement, like dYdX v3). This decision determines everything – regulatory load, tech stack, team.
Development proceeds in layers: first smart contracts with full Foundry coverage (fuzzing, invariant testing), then backend services, then integration layer, and finally frontend. Testing includes fork testing on mainnet via Foundry – we reproduce real liquidity conditions, not synthetic ones.
Audit is mandatory before mainnet deployment. For DEX contracts, minimally one firm with manual review (Trail of Bits, Spearbit, Code4rena contest). For CEX custody, audit of key storage processes. We guarantee all contracts undergo formal verification and fuzzing testing (Echidna, Foundry invariant).
Estimated timelines
| Exchange type |
Timeframe |
| DEX (AMM, xy=k) |
3 to 5 months |
| DEX with concentrated liquidity (v3-like) |
6 to 10 months |
| CEX (matching engine + custody + trading UI) |
8 to 14 months |
| Integration with existing protocol |
4 to 8 weeks |
Cost is calculated individually after a technical briefing: chain selection, throughput requirements, custodial model. Our certified engineers with 10+ years of experience will help you choose the optimal architecture and avoid common pitfalls. Contact our team for a detailed proposal.
Pitfalls to avoid at launch
- Forgetting the price oracle in AMM. Spot price can be manipulated with a flash loan in one transaction. If your lending protocol uses the spot price from its own pool, that's a bug.
- Hot wallet without limits. A CEX without daily limits on automatic withdrawals is an invitation for attackers. Compromising one key should lose at most 10% of total funds.
- Absence of circuit breaker. A 40% price drop in 5 minutes should halt automatic liquidations or withdrawals until manual review. Without this, a cascading liquidation spiral destroys all TVL.
- Incorrect decimal handling. USDC uses 6 decimals, WBTC – 8, most tokens – 18. Mixing without normalization leads to either precision loss or overflow. Solidity has no float; we work with fixed-point using FullMath (mulDiv with overflow protection).
Want to avoid these problems? Get a consultation — we will select the architecture for your project and provide exact timelines. Order exchange development with quality guarantee and ongoing support.