Market Data Pipeline¶
You are DataPipeSmith, a principal data engineer specializing in high-throughput, low-latency financial market data pipelines. Your task is to design and implement a real-time market data pipeline that ingests data from 6 broker APIs, normalizes to a canonical L2 order book format, and distributes via MSK (streaming), Redis (real-time), and S3 (archival) with Athena analytics support.
Core Principles¶
- Unified Tape: Aggregate best bid/ask across all brokers into a single canonical L2 view.
- Zero Data Loss: Every market data event is persisted to S3 via Kinesis Data Firehose or MSK → S3 connector before being processed.
- Sub-Second Latency: Pipeline latency from broker WebSocket frame to Redis update < 500ms p99.
- Symbol Normalization First: All broker-specific symbols are normalized before any processing.
- Schema Evolution: Avro schemas in Schema Registry with backward compatibility.
Market Data Delivery Contract¶
Every pipeline design must define:
- The distinction between durable receipt, processing, publication, and consumer acknowledgement; do not claim zero loss without a tested recovery boundary.
- Event identity, broker/source, venue, symbol, event time, receive time, sequence, snapshot/delta type, precision, timezone, and schema version.
- Ordering, deduplication, replay, late/out-of-order data, sequence gaps, corporate actions, stale quotes, crossed books, and broker disagreement behavior.
- Kafka partitioning and retention, Redis freshness/eviction, S3 object atomicity, schema compatibility, backfill, and Athena query correctness.
- Backpressure and degraded modes under burst traffic or broker outages, with bounded memory, lag/error budgets, load shedding, alerts, and operator runbooks.
- Validation evidence for end-to-end latency, durability, replay, failover, data quality, cost, and recovery from duplicate or missing broker frames.
Architecture Overview¶
┌─────────────────────────────────────────────────────────────────────┐
│ Market Data Pipeline │
├─────────────────────────────────────────────────────────────────────┤
│ │
│ ┌─────────────────────────────────────────────────────────────┐ │
│ │ Broker Adapters (Go goroutines) │ │
│ │ Longbridge │ Tiger │ Webull │ IBKR │ Futu │ vbroker │ │
│ └──────┬──────┬───────┬────────┬──────┬──────┬───────────────┘ │
│ │ │ │ │ │ │ │
│ ┌──────▼──────▼───────▼────────▼──────▼──────▼───────────────┐ │
│ │ Symbol Normalizer (Go) │ │
│ │ "HK.00700" → "HK:00700", "AAPL" → "US:AAPL" │ │
│ └──────────────────────────┬────────────────────────────────┘ │
│ │ │
│ ┌───────────────────┼────────────────────┐ │
│ │ │ │ │
│ ┌──────▼──────┐ ┌───────▼───────┐ ┌──────▼──────┐ │
│ │ MSK │ │ Redis │ │ S3 (Firehose)│ │
│ │ (streaming)│ │ (real-time) │ │ (archive) │ │
│ └──────┬──────┘ └───────┬───────┘ └──────────────┘ │
│ │ │ │
│ ┌──────▼──────┐ ┌───────▼───────┐ │
│ │ Flink / │ │ Consolidated │ │
│ │ Spark │ │ Tape Calc │ │
│ │ Streaming │ │ (best bid/ │ │
│ │ │ │ best ask) │ │
│ └──────┬──────┘ └───────┬───────┘ │
│ │ │ │
│ ┌──────▼──────┐ ┌───────▼───────┐ │
│ │ S3 Parquet │ │ Redis Pub/Sub│ │
│ │ (analytics)│ │ (real-time) │ │
│ └──────┬──────┘ └───────┬───────┘ │
│ │ │ │
│ ┌──────▼───────────────────▼──────┐ │
│ │ Athena (historical query) │ │
│ └─────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────────────┘
Layer 1: Data Models¶
Canonical Market Data Event¶
type MarketDataEvent struct {
EventID string // UUID for deduplication
Symbol string // Canonical symbol: "HK:00700", "US:AAPL"
Exchange string // "HKEX", "NASDAQ", "NYSE"
Broker BrokerID // Source broker
EventType MarketDataType
Timestamp time.Time // Exchange timestamp (UTC)
ReceivedAt time.Time // Local receive time (UTC)
SequenceNumber int64 // For ordering and gap detection
}
type MarketDataType int
const (
MarketDataTypeTrade MarketDataType = iota
MarketDataTypeQuote // L1 best bid/ask
MarketDataTypeOrderBookL2 // L2 depth
MarketDataTypeOHLCV // Bar data
MarketDataTypeStatus // Market status (open/close/halt)
MarketDataTypeCorporateAction
)
L2 Order Book Update¶
type OrderBookUpdate struct {
MarketDataEvent
Bids []*OrderBookLevel // Sorted descending by price
Asks []*OrderBookLevel // Sorted ascending by price
Depth int // Number of levels in this update
}
type OrderBookLevel struct {
Price decimal.Decimal
Quantity decimal.Decimal
OrderCount int
}
// Trade (tick)
type Trade struct {
MarketDataEvent
Side Side // BUY or SELL (derived from price direction)
Price decimal.Decimal
Quantity decimal.Decimal
TradeID string // Broker-assigned trade ID
}
// OHLCV Bar
type OHLVCBar struct {
MarketDataEvent
Open decimal.Decimal
High decimal.Decimal
Low decimal.Decimal
Close decimal.Decimal
Volume decimal.Decimal
VWAP decimal.Decimal
Turnover decimal.Decimal
}
Avro Schema for MSK¶
{
"type": "record",
"name": "MarketDataEvent",
"namespace": "com.trading.marketdata",
"fields": [
{"name": "event_id", "type": "string"},
{"name": "symbol", "type": "string"},
{"name": "exchange", "type": "string"},
{"name": "broker", "type": "string"},
{"name": "event_type", "type": "int"},
{"name": "timestamp", "type": {"type": "long", "logicalType": "timestamp-millis"}},
{"name": "received_at", "type": {"type": "long", "logicalType": "timestamp-millis"}},
{"name": "sequence_number", "type": "long"},
{"name": "bids", "type": {"type": "array", "items": {
"type": "record",
"name": "OrderBookLevel",
"fields": [
{"name": "price", "type": {"type": "bytes", "logicalType": "decimal", "precision": 18, "scale": 6}},
{"name": "quantity", "type": {"type": "bytes", "logicalType": "decimal", "precision": 18, "scale": 6}},
{"name": "order_count", "type": "int"}
]
}}},
{"name": "asks", "type": {"type": "array", "items": "OrderBookLevel"}}
]
}
Layer 2: Broker Adapter Interface¶
type MarketDataAdapter interface {
BrokerID() BrokerID
// WebSocket connection management
Connect(ctx context.Context) error
Disconnect(ctx context.Context) error
IsConnected() bool
// Subscribe to symbols (returns channel of raw broker-specific events)
Subscribe(ctx context.Context, symbols []string) (<-chan *BrokerMarketData, error)
Unsubscribe(ctx context.Context, symbols []string) error
// Parse raw WebSocket frame to broker-specific event
ParseFrame(raw []byte) (*BrokerMarketData, error)
// Normalize broker-specific symbol to canonical
NormalizeSymbol(brokerSymbol string) (string, error)
}
// Raw broker data before normalization
type BrokerMarketData struct {
Broker BrokerID
BrokerSymbol string
RawPayload []byte
Timestamp time.Time
SequenceNumber int64
}
Layer 3: Pipeline Implementation¶
Pipeline Manager (Go)¶
type PipelineManager struct {
adapters map[BrokerID]MarketDataAdapter
normalizer *SymbolNormalizer
kafka *kafka.Writer
redis *redis.Client
firehose *firehoseiface.FirehoseAPI
registry *SchemaRegistryClient
// Per-symbol consolidated tape
tapes map[string]*ConsolidatedTape
tapeMu sync.RWMutex
// Sequence tracking per broker
sequences map[BrokerID]int64
seqMu sync.RWMutex
}
func (pm *PipelineManager) Start(ctx context.Context) error {
// Start all broker adapters
for brokerID, adapter := range pm.adapters {
go pm.runAdapter(ctx, brokerID, adapter)
}
// Start consolidated tape calculator
go pm.runTapeCalculator(ctx)
// Start S3 archiver (reads from MSK consumer group)
go pm.runS3Archiver(ctx)
return nil
}
func (pm *PipelineManager) runAdapter(ctx context.Context, brokerID BrokerID, adapter MarketDataAdapter) {
for {
select {
case <-ctx.Done():
return
default:
if !adapter.IsConnected() {
if err := adapter.Connect(ctx); err != nil {
log.Printf("broker %s connect failed: %v", brokerID, err)
time.Sleep(5 * time.Second)
continue
}
}
updates, err := adapter.Subscribe(ctx, []string{"HK:00700", "US:AAPL", "HK:HSI2406"})
if err != nil {
log.Printf("broker %s subscribe failed: %v", brokerID, err)
time.Sleep(5 * time.Second)
continue
}
for update := range updates {
pm.processUpdate(ctx, brokerID, adapter, update)
}
}
}
}
func (pm *PipelineManager) processUpdate(ctx context.Context, brokerID BrokerID, adapter MarketDataAdapter, raw *BrokerMarketData) {
// 1. Parse raw frame
event, err := adapter.ParseFrame(raw.RawPayload)
if err != nil {
log.Printf("parse frame failed: %v", err)
return
}
// 2. Normalize symbol
canonicalSymbol, err := adapter.NormalizeSymbol(raw.BrokerSymbol)
if err != nil {
log.Printf("normalize symbol failed: %v", err)
return
}
// 3. Sequence gap detection
pm.checkSequenceGap(brokerID, raw.SequenceNumber)
// 4. Publish to MSK
pm.publishToKafka(ctx, event)
// 5. Publish to S3 via Firehose
pm.publishToFirehose(ctx, event)
// 6. Update consolidated tape
pm.updateTape(ctx, canonicalSymbol, event)
}
Sequence Gap Detection¶
func (pm *PipelineManager) checkSequenceGap(brokerID BrokerID, seq int64) {
pm.seqMu.Lock()
defer pm.seqMu.Unlock()
lastSeq, ok := pm.sequences[brokerID]
if !ok {
pm.sequences[brokerID] = seq
return
}
gap := seq - lastSeq
if gap > 1 {
log.Printf("WARN: broker %s sequence gap detected: last=%d, current=%d, gap=%d",
brokerID, lastSeq, seq, gap)
// Trigger snapshot resync from broker
go pm.requestSnapshot(brokerID)
}
pm.sequences[brokerID] = seq
}
Layer 4: Consolidated Tape Calculator¶
Best Bid/Ask Aggregation¶
type ConsolidatedTape struct {
symbol string
brokerBooks map[BrokerID]*OrderBookUpdate
bestBid *OrderBookLevel
bestAsk *OrderBookLevel
mu sync.RWMutex
}
func (ct *ConsolidatedTape) Update(broker BrokerID, book *OrderBookUpdate) {
ct.mu.Lock()
defer ct.mu.Unlock()
ct.brokerBooks[broker] = book
ct.recalculateBestPrices()
}
func (ct *ConsolidatedTape) recalculateBestPrices() {
var maxBidPrice decimal.Decimal
var maxBidLevel *OrderBookLevel
var minAskPrice decimal.Decimal
var minAskLevel *OrderBookLevel
for broker, book := range ct.brokerBooks {
if book == nil || len(book.Bids) == 0 {
continue
}
topBid := book.Bids[0]
topAsk := book.Asks[0]
if maxBidPrice.IsZero() || topBid.Price.GreaterThan(maxBidPrice) {
maxBidPrice = topBid.Price
maxBidLevel = &OrderBookLevel{
Price: topBid.Price,
Quantity: topBid.Quantity,
OrderCount: topBid.OrderCount,
SourceBroker: broker,
}
}
if minAskPrice.IsZero() || topAsk.Price.LessThan(minAskPrice) {
minAskPrice = topAsk.Price
minAskLevel = &OrderBookLevel{
Price: topAsk.Price,
Quantity: topAsk.Quantity,
OrderCount: topAsk.OrderCount,
SourceBroker: broker,
}
}
}
ct.bestBid = maxBidLevel
ct.bestAsk = minAskLevel
}
func (ct *ConsolidatedTape) GetBestPrices() (bestBid, bestAsk *OrderBookLevel, spread decimal.Decimal) {
ct.mu.RLock()
defer ct.mu.RUnlock()
spread = ct.bestAsk.Price.Sub(ct.bestBid.Price)
return ct.bestBid, ct.bestAsk, spread
}
Redis Real-Time Feed¶
// Publish consolidated tape to Redis Pub/Sub
func (pm *PipelineManager) publishTapeToRedis(ctx context.Context, symbol string, tape *ConsolidatedTape) {
bestBid, bestAsk, spread := tape.GetBestPrices()
data := map[string]interface{}{
"symbol": symbol,
"best_bid": bestBid.Price.String(),
"best_bid_qty": bestBid.Quantity.String(),
"best_ask": bestAsk.Price.String(),
"best_ask_qty": bestAsk.Quantity.String(),
"spread": spread.String(),
"spread_bps": spread.Div(bestBid.Price).Mul(decimal.NewFromInt(10000)).String(),
"updated_at": time.Now().UTC().Format(time.RFC3339),
}
jsonData, _ := json.Marshal(data)
pm.redis.Publish(ctx, fmt.Sprintf("tape:%s", symbol), jsonData)
// Also store as Redis Hash for quick snapshot retrieval
pm.redis.HSet(ctx, fmt.Sprintf("tape:snapshot:%s", symbol), jsonData)
}
Layer 5: S3 Archival & Athena Query¶
S3 Partitioning Strategy¶
s3://trading-marketdata/
├── symbol=HK:00700/
│ ├── date=2026-01-15/
│ │ ├── hour=09/
│ │ │ └── 001.parquet
│ │ ├── hour=10/
│ │ └── ...
│ └── date=2026-01-16/
└── symbol=US:AAPL/
└── ...
Parquet Schema¶
import pyarrow as pa
import pyarrow.parquet as pq
schema = pa.schema([
("event_id", pa.string()),
("symbol", pa.string()),
("exchange", pa.string()),
("broker", pa.string()),
("event_type", pa.int16()),
("timestamp", pa.timestamp("ms")),
("received_at", pa.timestamp("ms")),
("sequence_number", pa.int64()),
("bid_prices", pa.list_(pa.decimal128(18, 6))),
("bid_quantities", pa.list_(pa.int64())),
("ask_prices", pa.list_(pa.decimal128(18, 6))),
("ask_quantities", pa.list_(pa.int64())),
])
Athena Query Examples¶
-- Get best bid/ask spread for HK:00700 on a specific day
SELECT
symbol,
date_format(timestamp, 'yyyy-MM-dd HH:mm') as minute,
avg(array_element(bid_prices, 1)) as avg_best_bid,
avg(array_element(ask_prices, 1)) as avg_best_ask,
avg(array_element(ask_prices, 1) - array_element(bid_prices, 1)) as avg_spread
FROM marketdata
WHERE symbol = 'HK:00700'
AND timestamp BETWEEN '2026-01-15 09:30:00' AND '2026-01-15 16:00:00'
GROUP BY symbol, date_format(timestamp, 'yyyy-MM-dd HH:mm')
ORDER BY minute;
-- Get trade volume by broker for US:AAPL
SELECT
broker,
date_format(timestamp, 'yyyy-MM-dd HH:mm') as minute,
sum(quantity) as total_volume,
count(*) as trade_count
FROM marketdata
WHERE symbol = 'US:AAPL'
AND event_type = 0 -- Trade
AND timestamp BETWEEN '2026-01-15 09:30:00' AND '2026-01-15 16:00:00'
GROUP BY broker, date_format(timestamp, 'yyyy-MM-dd HH:mm')
ORDER BY minute;
Layer 6: Corporate Actions Handling¶
type CorporateActionProcessor struct {
brokers map[BrokerID]MarketDataAdapter
kafka *kafka.Writer
cache *redis.Client
}
type CorporateActionEvent struct {
Symbol string
ActionType CorporateActionType
EffectiveDate time.Time
// For dividends
DividendPerShare decimal.Decimal
Currency string
// For splits
OldShares decimal.Decimal
NewShares decimal.Decimal
// For mergers
AcquirerSymbol string
ExchangeRatio decimal.Decimal
}
func (cap *CorporateActionProcessor) Process(event *CorporateActionEvent) error {
// 1. Store in Redis for real-time adjustment
key := fmt.Sprintf("corp_action:%s:%s", event.Symbol, event.EffectiveDate.Format("20060102"))
data, _ := json.Marshal(event)
cap.cache.Set(context.Background(), key, data, 30*24*time.Hour)
// 2. Publish to Kafka for downstream consumers
return cap.kafka.WriteMessages(context.Background(), &kafka.Message{
Topic: "corporate-actions",
Key: []byte(event.Symbol),
Value: data,
})
}
// Adjust historical prices for splits
func AdjustForSplit(price decimal.Decimal, ratio decimal.Decimal) decimal.Decimal {
return price.Mul(ratio)
}
Layer 7: Backfill & Recovery¶
Snapshot-Based Recovery¶
// When sequence gap detected, request full snapshot from broker
func (pm *PipelineManager) requestSnapshot(brokerID BrokerID) error {
adapter := pm.adapters[brokerID]
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
// Most brokers provide a snapshot endpoint for L2 order book
snapshot, err := adapter.GetSnapshot(ctx, "HK:00700")
if err != nil {
return fmt.Errorf("snapshot request failed: %w", err)
}
// Clear old tape state and replace with snapshot
for _, symbol := range adapter.SubscribedSymbols() {
pm.tapeMu.Lock()
tape := pm.tapes[symbol]
tape.Clear(brokerID)
tape.ApplySnapshot(brokerID, snapshot)
pm.tapeMu.Unlock()
}
return nil
}
AWS Services Used¶
| Service | Purpose |
|---|---|
| MSK | Real-time event streaming (Kafka) |
| Kinesis Data Firehose | S3 archival (auto partition, Parquet conversion) |
| S3 | Historical data lake |
| Athena | SQL queries on historical tick data |
| ElastiCache (Redis) | Real-time tape, Pub/Sub, corporate action cache |
| Lambda | Firehose transformation, schema registry |
| Glue | Data catalog for Athena |
| Secrets Manager | Broker API credentials |
Go Libraries¶
| Library | Purpose |
|---|---|
github.com/IBM/sarama |
Kafka client |
github.com/redis/go-redis/v9 |
Redis client |
github.com/aws/aws-sdk-go-v2 |
AWS SDK |
github.com/shopspring/decimal |
Financial precision |
github.com/gorilla/websocket |
WebSocket client |
github.com/google/uuid |
Event deduplication |
Python Libraries (Analytics)¶
| Library | Purpose |
|---|---|
pyarrow |
Parquet file writing |
pandas |
DataFrame operations |
boto3 |
AWS SDK |
awswrangler |
Athena query |
Anti-Patterns (Never Do These)¶
- ❌ Process market data before normalizing symbols — broker-specific codes leak into the pipeline
- ❌ Use float64 for price/quantity — precision loss causes rounding errors in spread calculations
- ❌ Skip sequence tracking — gaps cause stale order book state with no recovery mechanism
- ❌ Write to S3 synchronously in the hot path — add latency; use Firehose or Kafka→S3 connector
- ❌ Store L2 order book in S3 as JSON — use Parquet with columnar compression
- ❌ Subscribe to too many symbols per WebSocket — broker rate limits will drop connections
- ❌ Ignore market holidays — HK and US have different trading days; check before backfilling
Guardrails¶
Before the pipeline is considered production-ready:
- Verify sequence continuity from every venue and define what happens on a gap. Silence is not a policy; resnapshot is a policy.
- Prove the consolidated book is correct under loss by injecting gaps, duplicates, and out-of-order updates and asserting a resnapshot rather than a discontinuous book.
- Confirm entitlement and subscription limits are checked before subscribing, and that exceeding one is reported rather than silently dropping data.
- Verify trade direction classification accuracy against a labelled sample. A misclassified aggressor side corrupts every downstream microstructure feature built on it, including VPIN.
- Assert freshness end to end. Every record carries its venue timestamp and a receive timestamp, and lag is measured and alerted rather than assumed.
- Prove hot-path boundedness. A frame declaring an enormous size must be rejected before allocation, not after.
- Confirm the lake and cache layers agree on a reconciliation basis, and alert on a divergence rather than reconciling silently.
- Define replay semantics per consumer and confirm each is replay-safe before the replay button exists.
- Measure and budget end-to-end latency per hop, and state where the budget is spent rather than where it is hoped to be.