The Pipeline Architecture
A trading data pipeline has four stages: collection, validation, storage, and serving. Each stage can fail independently, and each failure mode requires a different response.
[Collection] → [Validation] → [Storage] → [Serving]
| | | |
API calls Sanity checks PostgreSQL Redis cache
Scheduling Gap detection TimescaleDB REST API
Retry logic Deduplication Backfill Rate limits Stage 1: Collection
import requests
import schedule
import time
import logging
from datetime import datetime
logger = logging.getLogger("pipeline")
API_KEY = "YOUR_API_KEY"
BASE_URL = "https://tickatlas.com/v1"
SYMBOLS = ["EURUSD", "GBPUSD", "USDJPY", "XAUUSD", "BTCUSD"]
TIMEFRAMES = ["M15", "H1", "H4", "D1"]
INDICATORS = ["RSI_14", "MACD_hist", "EMA_50", "SMA_200", "ATR_14"]
def collect_snapshot():
"""Collect full indicator snapshot for all symbols."""
timestamp = datetime.utcnow()
records = []
for symbol in SYMBOLS:
for tf in TIMEFRAMES:
for indicator in INDICATORS:
try:
resp = requests.get(
f"{BASE_URL}/indicator",
headers={"X-API-Key": API_KEY},
params={
"symbol": symbol,
"indicator": indicator,
"timeframe": tf,
},
timeout=10,
)
resp.raise_for_status()
data = resp.json()["data"]
records.append({
"collected_at": timestamp.isoformat(),
"symbol": symbol,
"timeframe": tf,
"indicator": indicator,
"value": data["value"],
"bid": data.get("bid"),
"ask": data.get("ask"),
})
except Exception as e:
logger.error(f"Collection failed: {symbol}/{tf}/{indicator}: {e}")
logger.info(f"Collected {len(records)} records at {timestamp}")
return records
# Schedule collection
schedule.every(5).minutes.do(collect_snapshot) Stage 2: Validation
from dataclasses import dataclass
@dataclass
class ValidationResult:
valid: bool
errors: list[str]
def validate_record(record: dict) -> ValidationResult:
"""Validate a collected data record."""
errors = []
# Check required fields
for field in ["symbol", "timeframe", "indicator", "value"]:
if field not in record:
errors.append(f"Missing field: {field}")
# Validate the indicator value is numeric (the API returns a bare float or null)
val = record.get("value")
if not isinstance(val, (int, float)):
errors.append(f"Non-numeric value for {record.get('indicator')}: {val}")
elif val != val: # NaN check
errors.append(f"NaN value for {record.get('indicator')}")
# Validate quote sanity
bid, ask = record.get("bid"), record.get("ask")
if bid is not None and ask is not None:
if ask < bid:
errors.append("Ask is below bid")
if bid <= 0:
errors.append("Bid is zero or negative")
# Validate RSI range
if record.get("indicator") == "RSI_14" and isinstance(val, (int, float)):
if not (0 <= val <= 100):
errors.append(f"RSI out of range: {val}")
return ValidationResult(valid=len(errors) == 0, errors=errors)
# Usage
for record in records:
result = validate_record(record)
if not result.valid:
logger.warning(f"Invalid record: {record['symbol']}: {result.errors}") Stage 3: Storage
import psycopg2
import json
def store_records(records: list[dict], conn_string: str):
"""Batch insert validated records into PostgreSQL."""
conn = psycopg2.connect(conn_string)
cur = conn.cursor()
insert_sql = """
INSERT INTO indicator_snapshots
(collected_at, symbol, timeframe, indicator, value, bid, ask)
VALUES (%s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (symbol, timeframe, indicator, collected_at)
DO NOTHING
"""
batch = []
for record in records:
result = validate_record(record)
if result.valid:
batch.append((
record["collected_at"],
record["symbol"],
record["timeframe"],
record["indicator"],
record["value"],
record.get("bid"),
record.get("ask"),
))
cur.executemany(insert_sql, batch)
conn.commit()
logger.info(f"Stored {len(batch)} records ({len(records) - len(batch)} invalid)")
cur.close()
conn.close() Stage 4: Serving via Cache
import redis
r = redis.Redis(decode_responses=True)
def update_cache(records: list[dict]):
"""Push latest values to Redis for fast access."""
pipe = r.pipeline()
for record in records:
key = f"latest:{record['symbol']}:{record['timeframe']}:{record['indicator']}"
pipe.setex(key, 600, json.dumps({
"value": record["value"],
"bid": record.get("bid"),
"ask": record.get("ask"),
"collected_at": record["collected_at"],
}))
pipe.execute()
def get_latest(symbol: str, timeframe: str, indicator: str) -> dict:
"""Fast retrieval from cache."""
key = f"latest:{symbol}:{timeframe}:{indicator}"
data = r.get(key)
return json.loads(data) if data else None Gap Detection
def detect_gaps(symbol: str, timeframe: str, expected_interval_minutes: int):
"""Check for missing data points in the collection history."""
conn = psycopg2.connect(CONN_STRING)
cur = conn.cursor()
cur.execute("""
SELECT collected_at FROM indicator_snapshots
WHERE symbol = %s AND timeframe = %s AND indicator = 'RSI_14'
ORDER BY collected_at DESC LIMIT 100
""", (symbol, timeframe))
timestamps = [row[0] for row in cur.fetchall()]
gaps = []
for i in range(len(timestamps) - 1):
diff = (timestamps[i] - timestamps[i + 1]).total_seconds() / 60
if diff > expected_interval_minutes * 2:
gaps.append({
"from": timestamps[i + 1].isoformat(),
"to": timestamps[i].isoformat(),
"gap_minutes": int(diff),
})
cur.close()
conn.close()
return gaps Run this against live data.
Every account starts pay-as-you-go with $2.50 of credit and no card. Paste the key into the samples above and the requests work unchanged.