Real-World Project
Financial Time-Series Analytics
Build a high-performance stock market data platform with Cassandra!
๐ Project Overview
Build a real-world financial platform!
The Challenge: FinTrack Pro
FinTrack Pro provides real-time stock market data to 100,000+ traders worldwide. They need to:
- ๐น Track 10,000 stocks: Real-time tick data
- ๐ OHLC candlesticks: 1-min, 5-min, 15-min, 1-hour, 1-day
- ๐ Technical indicators: Moving averages, volume, etc
- ๐ Price alerts: Notify on price movements
- ๐ฑ Historical data: 10 years of daily data
Scale Requirements
| Stocks Tracked: | 10,000 symbols (NYSE, NASDAQ, etc) |
| Tick Data/Second: | ~50,000 writes/second (peak market hours) |
| Data Points/Day: | 1.8 billion ticks (6.5 hours ร 50K/sec) |
| Retention: | Ticks: 7 days, Candles: 10 years |
| Query Pattern: | Real-time charts, historical analysis |
| Users: | 100,000 concurrent traders |
Why Cassandra?
- โ High write throughput: 50K+ writes/second
- โ Time-series optimized: TWCS + clustering by time
- โ Fast range queries: Get candles for chart
- โ Multi-resolution: Different time buckets
- โ No downtime: 24/7 trading data availability
๐ฏ Requirements & Use Cases
What we need to build!
Data Ingestion
- Ingest tick data (50K/sec)
- Calculate OHLC candles
- Store multiple timeframes
- Handle late data arrivals
- Batch processing pipelines
Chart Queries
- Get candles for date range
- Real-time price updates
- Volume analysis
- Technical indicators
- Multiple resolution charts
Analytics
- Moving averages (SMA, EMA)
- Price change % calculations
- Volume-weighted averages
- Historical comparisons
- Statistical analysis
Alerts & Triggers
- Price threshold alerts
- Volume spike detection
- Technical indicator signals
- Real-time notifications
- Watchlist monitoring
๐บ๏ธ Data Model Design
Financial time-series modeling!
Time-Series Data Modeling Principles
- Bucketing: Partition by symbol + time window (day/week/month)
- Clustering: Order by timestamp DESC (recent first)
- Multi-resolution: Separate tables for each timeframe
- Denormalization: Store calculated values (open, close, high, low)
- TTL strategy: Different retention per resolution
Query Patterns โ Table Design
| Query Pattern | Access Pattern | Table Design |
|---|---|---|
| Q1: Real-time tick data | By symbol + timestamp | ticks_by_symbol_day |
| Q2: 1-minute candles | By symbol + time range | candles_1min |
| Q3: Daily OHLC | By symbol + date | candles_daily |
| Q4: Latest price | By symbol (most recent) | latest_prices |
| Q5: Historical range | By symbol + date range | candles_daily |
Data Flow Architecture
Market Data Feed
โ
โผ
โโโโโโโโโโโโโโโโ
โ Tick Ingesterโ โ 50K writes/sec
โโโโโโโโฌโโโโโโโโ
โ
โผ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ Cassandra Cluster โ
โ โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ โ
โ โ ticks_by_symbol_day (7 days) โ โ
โ โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ โ
โโโโโโโโโโโโโโโโฌโโโโโโโโโโโโโโโโโโโโโโโโ
โ
โผ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ Aggregation Pipeline (Spark) โ
โ - Group by 1-min windows โ
โ - Calculate OHLC, Volume โ
โ - Compute indicators โ
โโโโโโโโโโโโโโโโฌโโโโโโโโโโโโโโโโโโโโโโโโ
โ
โโโโโโโโโโโโโโผโโโโโโโโโโโโโ
โ โ โ
โผ โผ โผ
candles_1min candles_5min candles_daily
(30 days) (90 days) (10 years)
๐ Schema Design
Complete table schemas!
1
Keyspace Creation
-- Create keyspace
CREATE KEYSPACE fintrack_data
WITH replication = {
'class': 'NetworkTopologyStrategy',
'datacenter1': 3,
'datacenter2': 3
};
USE fintrack_data;
2
Tick Data Table (Raw Market Data)
-- Raw tick data partitioned by symbol + day
CREATE TABLE ticks_by_symbol_day (
symbol TEXT, -- Stock symbol (AAPL, GOOGL)
day DATE, -- Trading day (partition bucket)
tick_time TIMESTAMP, -- Exact time of trade
price DECIMAL, -- Trade price
volume BIGINT, -- Trade volume
bid DECIMAL, -- Bid price
ask DECIMAL, -- Ask price
exchange TEXT, -- Exchange (NYSE, NASDAQ)
PRIMARY KEY ((symbol, day), tick_time)
) WITH CLUSTERING ORDER BY (tick_time DESC)
AND default_time_to_live = 604800 -- 7 days
AND compaction = {
'class': 'TimeWindowCompactionStrategy',
'compaction_window_unit': 'DAYS',
'compaction_window_size': 1
};
-- Design notes:
-- โ
Partition: (symbol, day) - one partition per stock per day
-- โ
Clustering: tick_time DESC - newest ticks first
-- โ
TTL: 7 days - ticks only for recent analysis
-- โ
Bucket size: ~280K ticks/day/stock (manageable)
3
1-Minute Candlestick Table
-- 1-minute OHLC candles
CREATE TABLE candles_1min (
symbol TEXT,
week DATE, -- Week bucket (partition)
candle_time TIMESTAMP, -- Start of 1-min window
open DECIMAL,
high DECIMAL,
low DECIMAL,
close DECIMAL,
volume BIGINT,
trades INT, -- Number of trades in candle
vwap DECIMAL, -- Volume-weighted average price
PRIMARY KEY ((symbol, week), candle_time)
) WITH CLUSTERING ORDER BY (candle_time DESC)
AND default_time_to_live = 2592000 -- 30 days
AND compaction = {
'class': 'TimeWindowCompactionStrategy',
'compaction_window_unit': 'DAYS',
'compaction_window_size': 7
};
-- Design notes:
-- โ
Partition: (symbol, week) - ~2K candles/partition
-- โ
Week bucketing - prevents partition from growing too large
-- โ
Pre-calculated OHLC - fast chart rendering
4
Daily Candlestick Table
-- Daily OHLC candles (10 years history)
CREATE TABLE candles_daily (
symbol TEXT,
year INT, -- Year bucket (partition)
candle_date DATE, -- Trading date
open DECIMAL,
high DECIMAL,
low DECIMAL,
close DECIMAL,
volume BIGINT,
adjusted_close DECIMAL, -- Adjusted for splits/dividends
PRIMARY KEY ((symbol, year), candle_date)
) WITH CLUSTERING ORDER BY (candle_date DESC);
-- No TTL - keep 10 years of daily data!
-- Partition: ~252 trading days/year = manageable size
5
Latest Prices Table (Real-Time)
-- Current price snapshot (fast lookup)
CREATE TABLE latest_prices (
symbol TEXT PRIMARY KEY,
last_price DECIMAL,
last_update TIMESTAMP,
day_open DECIMAL,
day_high DECIMAL,
day_low DECIMAL,
day_volume BIGINT,
prev_close DECIMAL,
change_percent DECIMAL
);
-- Simple table - one row per symbol
-- Updated in real-time for dashboard
6
Price Alerts Table
-- User-configured price alerts
CREATE TABLE price_alerts_by_user (
user_id UUID,
alert_id TIMEUUID,
symbol TEXT,
alert_type TEXT, -- 'above', 'below', 'change%'
threshold DECIMAL,
triggered BOOLEAN,
created_at TIMESTAMP,
PRIMARY KEY (user_id, alert_id)
);
-- Also need alerts by symbol for efficient checking
CREATE TABLE price_alerts_by_symbol (
symbol TEXT,
alert_id TIMEUUID,
user_id UUID,
alert_type TEXT,
threshold DECIMAL,
PRIMARY KEY (symbol, alert_id)
);
๐ป Implementation
Python implementation!
Tick Data Ingestion
# tick_ingestion.py
from cassandra.cluster import Cluster
from cassandra.query import BatchStatement
from datetime import datetime, date
from decimal import Decimal
class TickIngestion:
def __init__(self, contact_points):
cluster = Cluster(contact_points)
self.session = cluster.connect('fintrack_data')
# Prepared statements
self.insert_tick = self.session.prepare("""
INSERT INTO ticks_by_symbol_day
(symbol, day, tick_time, price, volume, bid, ask, exchange)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
""")
self.update_latest = self.session.prepare("""
UPDATE latest_prices
SET last_price = ?,
last_update = ?,
day_high = ?,
day_low = ?,
day_volume = day_volume + ?
WHERE symbol = ?
""")
def ingest_tick(self, tick):
"""Insert single tick and update latest price"""
symbol = tick['symbol']
tick_time = tick['timestamp']
today = tick_time.date()
# Insert tick data
self.session.execute(
self.insert_tick,
(
symbol,
today,
tick_time,
Decimal(str(tick['price'])),
tick['volume'],
Decimal(str(tick['bid'])),
Decimal(str(tick['ask'])),
tick['exchange']
)
)
# Update latest price (for real-time dashboard)
self.session.execute(
self.update_latest,
(
Decimal(str(tick['price'])),
tick_time,
Decimal(str(tick['price'])), # Update high
Decimal(str(tick['price'])), # Update low
tick['volume'],
symbol
)
)
def ingest_batch(self, ticks):
"""Batch insert for high throughput"""
batch = BatchStatement()
for tick in ticks:
today = tick['timestamp'].date()
batch.add(
self.insert_tick,
(
tick['symbol'],
today,
tick['timestamp'],
Decimal(str(tick['price'])),
tick['volume'],
Decimal(str(tick['bid'])),
Decimal(str(tick['ask'])),
tick['exchange']
)
)
self.session.execute(batch)
# Usage
ingestion = TickIngestion(['localhost'])
# Real-time tick from market feed
tick = {
'symbol': 'AAPL',
'timestamp': datetime.now(),
'price': 178.25,
'volume': 100,
'bid': 178.24,
'ask': 178.26,
'exchange': 'NASDAQ'
}
ingestion.ingest_tick(tick)
Query Service
# market_queries.py
from datetime import datetime, timedelta
class MarketQueries:
def __init__(self, contact_points):
cluster = Cluster(contact_points)
self.session = cluster.connect('fintrack_data')
def get_latest_price(self, symbol):
"""Get current price"""
query = "SELECT * FROM latest_prices WHERE symbol = ?"
row = self.session.execute(query, (symbol,)).one()
if row:
return {
'symbol': row.symbol,
'price': float(row.last_price),
'change_percent': float(row.change_percent),
'volume': row.day_volume,
'high': float(row.day_high),
'low': float(row.day_low)
}
return None
def get_intraday_candles(self, symbol, minutes=60):
"""Get 1-minute candles for last N minutes"""
now = datetime.now()
start_time = now - timedelta(minutes=minutes)
# Calculate week bucket
week_start = now - timedelta(days=now.weekday())
week = week_start.date()
query = """
SELECT * FROM candles_1min
WHERE symbol = ? AND week = ?
AND candle_time >= ? AND candle_time <= ?
"""
rows = self.session.execute(
query,
(symbol, week, start_time, now)
)
return [
{
'time': row.candle_time,
'open': float(row.open),
'high': float(row.high),
'low': float(row.low),
'close': float(row.close),
'volume': row.volume
}
for row in rows
]
def get_daily_candles(self, symbol, start_date, end_date):
"""Get daily candles for date range"""
# Query across multiple year partitions if needed
years = range(start_date.year, end_date.year + 1)
all_candles = []
for year in years:
query = """
SELECT * FROM candles_daily
WHERE symbol = ? AND year = ?
AND candle_date >= ? AND candle_date <= ?
"""
rows = self.session.execute(
query,
(symbol, year, start_date, end_date)
)
for row in rows:
all_candles.append({
'date': row.candle_date,
'open': float(row.open),
'high': float(row.high),
'low': float(row.low),
'close': float(row.close),
'volume': row.volume
})
return sorted(all_candles, key=lambda x: x['date'])
# Usage
queries = MarketQueries(['localhost'])
# Latest price
price = queries.get_latest_price('AAPL')
print(f"AAPL: ${price['price']} ({price['change_percent']}%)")
# Intraday chart (last hour)
candles = queries.get_intraday_candles('AAPL', minutes=60)
๐ Candlestick Aggregation
Convert ticks to OHLC candles!
# candle_aggregator.py - Real-time aggregation
from datetime import datetime, timedelta
from decimal import Decimal
class CandleAggregator:
def __init__(self, session):
self.session = session
self.insert_1min = session.prepare("""
INSERT INTO candles_1min
(symbol, week, candle_time, open, high, low, close, volume, trades, vwap)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""")
def aggregate_1min_candle(self, symbol, candle_start):
"""Aggregate 1-minute candle from tick data"""
candle_end = candle_start + timedelta(minutes=1)
day = candle_start.date()
# Get all ticks in 1-minute window
query = """
SELECT price, volume FROM ticks_by_symbol_day
WHERE symbol = ? AND day = ?
AND tick_time >= ? AND tick_time < ?
"""
rows = self.session.execute(
query,
(symbol, day, candle_start, candle_end)
)
ticks = list(rows)
if not ticks:
return
# Calculate OHLC
prices = [float(t.price) for t in ticks]
volumes = [t.volume for t in ticks]
ohlc = {
'open': prices[0], # First tick
'high': max(prices),
'low': min(prices),
'close': prices[-1], # Last tick
'volume': sum(volumes),
'trades': len(ticks)
}
# Calculate VWAP (Volume-Weighted Average Price)
total_value = sum(p * v for p, v in zip(prices, volumes))
ohlc['vwap'] = total_value / ohlc['volume'] if ohlc['volume'] > 0 else 0
# Calculate week bucket
week_start = candle_start - timedelta(days=candle_start.weekday())
week = week_start.date()
# Insert candle
self.session.execute(
self.insert_1min,
(
symbol,
week,
candle_start,
Decimal(str(ohlc['open'])),
Decimal(str(ohlc['high'])),
Decimal(str(ohlc['low'])),
Decimal(str(ohlc['close'])),
ohlc['volume'],
ohlc['trades'],
Decimal(str(ohlc['vwap']))
)
)
return ohlc
def run_aggregation_job(self, symbols):
"""Run every minute to aggregate previous minute"""
now = datetime.now()
# Round down to last complete minute
candle_start = now.replace(second=0, microsecond=0) - timedelta(minutes=1)
for symbol in symbols:
self.aggregate_1min_candle(symbol, candle_start)
print(f"Aggregated {len(symbols)} symbols for {candle_start}")
# Schedule to run every minute
import schedule
aggregator = CandleAggregator(session)
symbols = ['AAPL', 'GOOGL', 'MSFT', 'TSLA', ...]
schedule.every().minute.at(":05").do(aggregator.run_aggregation_job, symbols)
โก Performance Optimization
Production optimization!
Bucketing Strategy
- Ticks: By symbol + day (~280K/partition)
- 1-min candles: By symbol + week (~2K/partition)
- Daily candles: By symbol + year (~252/partition)
- Why? Keep partitions < 100MB
- Benefit: Fast queries, efficient compaction
Storage Optimization
- TTL: Ticks 7d, 1-min 30d, daily forever
- TWCS: Time-window compaction
- Compression: LZ4 (3:1 ratio)
- Bloom filters: Reduce reads
- Result: 1TB/day โ 330GB compressed
Query Optimization
- Prepared statements: 10x faster
- Latest prices table: No range scan
- Batch inserts: 100 ticks/batch
- Async writes: Non-blocking
- Read consistency: ONE (fast)
Caching Layer
- Redis: Latest prices (sub-second)
- CDN: Static historical candles
- Application cache: Hot stocks
- TTL: 1-5 seconds for real-time data
- Hit rate: 95%+ on popular stocks
๐ Deployment Architecture
Production deployment!
Production Architecture
Market Data Providers
(NASDAQ, NYSE)
โ
โผ
โโโโโโโโโโโโโโโโโโ
โ Load Balancer โ
โโโโโโโโโโฌโโโโโโโโ
โ
โโโโโโโโโโโโโโผโโโโโโโโโโโโโ
โ โ โ
โผ โผ โผ
โโโโโโโโโโโโโ โโโโโโโโโโโโโ โโโโโโโโโโโโโ
โ Ingestor 1โ โ Ingestor 2โ โ Ingestor 3โ
โ (50K/sec) โ โ (50K/sec) โ โ (50K/sec) โ
โโโโโโโฌโโโโโโ โโโโโโโฌโโโโโโ โโโโโโโฌโโโโโโ
โ โ โ
โโโโโโโโโโโโโโโผโโโโโโโโโโโโโโ
โ
โโโโโโโโโโโโโดโโโโโโโโโโโโ
โ Cassandra Cluster โ
โ โโโโโโโโโโโโโโโโโโโ โ
โ โ DC1: 6 nodes โ โ
โ โ DC2: 6 nodes โ โ
โ โโโโโโโโโโโโโโโโโโโ โ
โโโโโโโโโโโโโฌโโโโโโโโโโโโ
โ
โโโโโโโโโโโโโโโโโผโโโโโโโโโโโโโโโโ
โ โ โ
โผ โผ โผ
โโโโโโโโโโโโ โโโโโโโโโโโโ โโโโโโโโโโโโ
โ Redis โ โ Spark โ โ Query โ
โ (Cache) โ โ(Candles) โ โ API โ
โโโโโโโโโโโโ โโโโโโโโโโโโ โโโโโโโฌโโโโโ
โ
โโโโโโโโดโโโโโโโ
โ Clients โ
โ 100K users โ
โโโโโโโโโโโโโโโ
Production Checklist
| Component | Configuration | Status |
|---|---|---|
| Cluster Size | 12 nodes (6 per DC) | โ |
| Replication | RF=3 per DC | โ |
| Write Throughput | 50K writes/second | โ |
| Read Latency | p99 < 50ms | โ |
| Compaction | TWCS (1-day windows) | โ |
| Caching | Redis + CDN | โ |
| Monitoring | Prometheus + Grafana | โ |
| Backup | Daily snapshots to S3 | โ |
๐ Project Complete!
You built a production-ready financial time-series analytics platform!
๐ What You Built:
- ๐ Scale: 50K writes/sec, 1.8B ticks/day, 10K stocks
- ๐บ๏ธ Data model: Multi-resolution (ticks, 1-min, daily)
- ๐ Schema: 6 tables with smart bucketing (day/week/year)
- ๐ป Implementation: Tick ingestion + candle aggregation
- ๐ OHLC candles: Real-time 1-min + historical daily
- โก Optimization: TWCS, TTL, caching, bucketing
- ๐ Deployment: Multi-DC with 12 nodes
๐ก Key Learnings:
- Bucketing strategy: Partition by symbol + time window
- Multi-resolution: Different tables per timeframe
- Pre-calculated OHLC: Store open/high/low/close
- Smart TTL: Ticks 7d, candles 30d-10y
- Latest prices table: Fast lookup, no range scan
- Aggregation pipeline: Convert ticks โ candles
- Caching layer: Redis for sub-second latency
๐ Architecture Highlights:
- โ 6 tables: ticks, candles_1min/daily, latest_prices, alerts
- โ Bucketing: Day (ticks), Week (1-min), Year (daily)
- โ Write path: 50K/sec with batch + async
- โ Read path: Single-partition queries (p99 < 50ms)
- โ Aggregation: Spark job every minute for candles
- โ Caching: Redis (latest) + CDN (historical)
๐ Next Steps:
- ๐ Add technical indicators (SMA, EMA, RSI, MACD)
- ๐ Build real-time alert notification system
- ๐ฑ Create mobile app for traders
- ๐ค Add ML for price prediction
- ๐ Build advanced charting with TradingView
- ๐ Expand to crypto, forex, commodities
๐ฏ You're ready to build financial platforms! ๐
Apply this to ANY time-series use case!
Advertisement
Responsive Ad