""" DuckDB Tick Data Loader — Parquet → DuckDB for fast analytical queries. Converts the raw Parquet store (gzip-compressed JSON payloads) into a normalized DuckDB database with tables for L2 snapshots, trades, funding rates, and pre-computed microstructural rollups. Usage: python data/duckdb_load.py --data-dir data/raw --db data/normalized/ftdt_tick.db python data/duckdb_load.py --coin BTC --days 7 python data/duckdb_load.py --incremental # only load new data since last run Tables created: l2_snapshots — full book state at each update time trades — aggressor-side classified trades funding — funding rate history l2_rollup_1s — pre-computed 1s microprice/OFI/VPIN rollups """ from __future__ import annotations import argparse import gzip import json import logging import os import sys import time from datetime import date, datetime, timedelta, timezone from pathlib import Path from typing import Optional sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) import numpy as np logger = logging.getLogger(__name__) # Columns extracted from L2 snapshots L2_SNAPSHOT_COLS = [ "exchange_ts_ms", "local_ts", "coin", "best_bid", "best_ask", "mid_price", "microprice", "obi", "spread_bps", "bid_depth_10", "ask_depth_10", "bid_levels", "ask_levels", ] # Columns extracted from trades TRADE_COLS = [ "exchange_ts_ms", "local_ts", "coin", "price", "size", "side", "aggressor", ] class DuckDBLoader: """Load raw Parquet data into a DuckDB database.""" def __init__(self, db_path: str = "data/normalized/ftdt_tick.db"): self._db_path = Path(db_path) self._db_path.parent.mkdir(parents=True, exist_ok=True) self._conn = None @property def conn(self): if self._conn is None: try: import duckdb self._conn = duckdb.connect(str(self._db_path)) logger.info("Connected to DuckDB: %s", self._db_path) except ImportError: raise ImportError( "duckdb not installed. Run: pip install duckdb" ) return self._conn def init_schema(self): """Create tables if they don't exist.""" self.conn.execute(""" CREATE TABLE IF NOT EXISTS l2_snapshots ( exchange_ts_ms BIGINT NOT NULL, local_ts DOUBLE NOT NULL, coin VARCHAR NOT NULL, best_bid DOUBLE NOT NULL, best_ask DOUBLE NOT NULL, mid_price DOUBLE, microprice DOUBLE, obi DOUBLE, spread_bps DOUBLE, bid_depth_10 DOUBLE, ask_depth_10 DOUBLE, bid_levels INTEGER, ask_levels INTEGER, PRIMARY KEY (coin, exchange_ts_ms) ) """) self.conn.execute(""" CREATE TABLE IF NOT EXISTS trades ( exchange_ts_ms BIGINT NOT NULL, local_ts DOUBLE NOT NULL, coin VARCHAR NOT NULL, price DOUBLE NOT NULL, size DOUBLE NOT NULL, side VARCHAR, aggressor VARCHAR, ) """) self.conn.execute(""" CREATE TABLE IF NOT EXISTS funding ( exchange_ts_ms BIGINT NOT NULL, local_ts DOUBLE NOT NULL, coin VARCHAR NOT NULL, funding_rate DOUBLE NOT NULL, mark_px DOUBLE, annual_apr DOUBLE, PRIMARY KEY (coin, exchange_ts_ms) ) """) self.conn.execute(""" CREATE TABLE IF NOT EXISTS load_state ( coin VARCHAR PRIMARY KEY, channel VARCHAR, last_loaded_ts BIGINT, loaded_at TIMESTAMP DEFAULT now() ) """) logger.info("Schema initialized") def load_l2( self, data_dir: str, coin: str, start_date: str, end_date: str, ) -> int: """Load L2 book data from Parquet into DuckDB. Returns row count.""" from data.store import read_range msgs = read_range(data_dir, "l2book", coin.upper(), start_date, end_date) if not msgs: return 0 rows = [] for msg in msgs: payload = msg.get("payload", {}) levels = payload.get("levels", []) msg_type = payload.get("type", "snapshot") bids_dict = {} asks_dict = {} if msg_type == "snapshot" and isinstance(levels, list): if len(levels) >= 1: for bid in levels[0]: sz = float(bid.get("sz", 0)) if sz > 0: bids_dict[float(bid["px"])] = sz if len(levels) >= 2: for ask in levels[1]: sz = float(ask.get("sz", 0)) if sz > 0: asks_dict[float(ask["px"])] = sz if not bids_dict or not asks_dict: continue bid_prices = sorted(bids_dict.keys(), reverse=True) ask_prices = sorted(asks_dict.keys()) best_bid = bid_prices[0] best_ask = ask_prices[0] mid = (best_bid + best_ask) / 2.0 bid_depth_10 = sum(bids_dict[px] for px in bid_prices[:10]) ask_depth_10 = sum(asks_dict[px] for px in ask_prices[:10]) total_depth = bid_depth_10 + ask_depth_10 obi = (bid_depth_10 - ask_depth_10) / total_depth if total_depth > 0 else 0.0 w = bid_depth_10 / total_depth if total_depth > 0 else 0.5 microprice = w * best_bid + (1 - w) * best_ask spread_bps = (best_ask - best_bid) / mid * 10000 if mid > 0 else 0 rows.append(( msg.get("exchange_ts", 0) or 0, msg.get("local_ts", 0.0), coin.upper(), best_bid, best_ask, mid, microprice, obi, spread_bps, bid_depth_10, ask_depth_10, len(bid_prices), len(ask_prices), )) if rows: import duckdb rel = duckdb.from_sequence(rows) self.conn.execute( "INSERT OR IGNORE INTO l2_snapshots SELECT * FROM rel" ) logger.info("Loaded %d L2 snapshots for %s", len(rows), coin) return len(rows) def load_trades( self, data_dir: str, coin: str, start_date: str, end_date: str, ) -> int: """Load trade data from Parquet into DuckDB. Returns row count.""" from data.store import read_range msgs = read_range(data_dir, "trades", coin.upper(), start_date, end_date) if not msgs: return 0 rows = [] for msg in msgs: payload = msg.get("payload", {}) px = float(payload.get("px", 0)) sz = float(payload.get("sz", 0)) if px <= 0 or sz <= 0: continue side = str(payload.get("side", "?")) aggressor = "buy" if side.upper() in ("B", "BUY") else "sell" rows.append(( msg.get("exchange_ts", 0) or 0, msg.get("local_ts", 0.0), coin.upper(), px, sz, side, aggressor, )) if rows: import duckdb rel = duckdb.from_sequence(rows) self.conn.execute( "INSERT INTO trades SELECT * FROM rel" ) logger.info("Loaded %d trades for %s", len(rows), coin) return len(rows) def load_funding( self, data_dir: str, coin: str, start_date: str, end_date: str, ) -> int: """Load funding rate data from Parquet into DuckDB.""" from data.store import read_range msgs = read_range(data_dir, "funding", coin.upper(), start_date, end_date) if not msgs: return 0 rows = [] for msg in msgs: payload = msg.get("payload", {}) rate = float(payload.get("funding", 0)) mark = float(payload.get("mark_px", 0)) annual = rate * 1095 if rate else 0 rows.append(( msg.get("exchange_ts", 0) or 0, msg.get("local_ts", 0.0), coin.upper(), rate, mark, annual, )) if rows: import duckdb rel = duckdb.from_sequence(rows) self.conn.execute( "INSERT OR IGNORE INTO funding SELECT * FROM rel" ) logger.info("Loaded %d funding observations for %s", len(rows), coin) return len(rows) def create_rollups(self): """Create pre-computed 1-second and 1-minute aggregation views.""" self.conn.execute(""" CREATE VIEW IF NOT EXISTS l2_rollup_1s AS SELECT (exchange_ts_ms / 1000)::BIGINT * 1000 AS ts_1s, coin, AVG(mid_price) AS mid_price, AVG(microprice) AS microprice, AVG(obi) AS obi, AVG(spread_bps) AS spread_bps, AVG(bid_depth_10) AS bid_depth_10, AVG(ask_depth_10) AS ask_depth_10, COUNT(*) AS n_snapshots FROM l2_snapshots GROUP BY ts_1s, coin """) self.conn.execute(""" CREATE VIEW IF NOT EXISTS ofi_rollup_1s AS SELECT (exchange_ts_ms / 1000)::BIGINT * 1000 AS ts_1s, coin, SUM(CASE WHEN aggressor = 'buy' THEN size ELSE 0 END) AS buy_volume, SUM(CASE WHEN aggressor = 'sell' THEN size ELSE 0 END) AS sell_volume, COUNT(*) AS trade_count FROM trades GROUP BY ts_1s, coin """) self.conn.execute(""" CREATE VIEW IF NOT EXISTS micro_rollup_1s AS SELECT b.ts_1s, b.coin, b.mid_price, b.microprice, b.spread_bps, b.n_snapshots, COALESCE(o.buy_volume, 0) AS buy_volume, COALESCE(o.sell_volume, 0) AS sell_volume, COALESCE(o.trade_count, 0) AS trade_count, CASE WHEN COALESCE(o.buy_volume + o.sell_volume, 0) > 0 THEN (o.buy_volume - o.sell_volume)::DOUBLE / (o.buy_volume + o.sell_volume) ELSE 0.0 END AS trade_imbalance FROM l2_rollup_1s b LEFT JOIN ofi_rollup_1s o ON b.ts_1s = o.ts_1s AND b.coin = o.coin """) logger.info("Rollup views created") def load_all( self, data_dir: str, coins: list[str], start_date: str, end_date: str, ) -> dict: """Load all data for given coins and date range.""" self.init_schema() totals = {"l2": 0, "trades": 0, "funding": 0} for coin in coins: totals["l2"] += self.load_l2(data_dir, coin, start_date, end_date) totals["trades"] += self.load_trades(data_dir, coin, start_date, end_date) totals["funding"] += self.load_funding(data_dir, coin, start_date, end_date) self.create_rollups() self.conn.execute( "INSERT OR REPLACE INTO load_state (coin, channel, last_loaded_ts, loaded_at) " "VALUES ('ALL', 'all', ?::BIGINT, now())", [int(time.time() * 1000)], ) return totals def query_markouts( self, coin: str, start_date: str, end_date: str, horizons_ms: list[int] | None = None, ) -> dict: """Compute trade markouts directly in DuckDB.""" if horizons_ms is None: horizons_ms = [100, 500, 1000, 5000, 10000, 30000, 60000] start_ts = int(datetime.fromisoformat(start_date).timestamp() * 1000) end_ts = int(datetime.fromisoformat(end_date).timestamp() * 1000) query = """ WITH trade_mids AS ( SELECT t.exchange_ts_ms, t.price AS trade_px, t.size AS trade_sz, t.aggressor, t.coin, s.mid_price AS mid_at_trade FROM trades t LEFT JOIN l2_snapshots s ON s.coin = t.coin AND s.exchange_ts_ms <= t.exchange_ts_ms AND s.exchange_ts_ms >= t.exchange_ts_ms - 2000 WHERE t.coin = ?::VARCHAR AND t.exchange_ts_ms >= ?::BIGINT AND t.exchange_ts_ms < ?::BIGINT QUALIFY ROW_NUMBER() OVER ( PARTITION BY t.exchange_ts_ms, t.price, t.size ORDER BY ABS(s.exchange_ts_ms - t.exchange_ts_ms) ) = 1 ) SELECT aggressor, COUNT(*) AS n, AVG(markout_bps) AS mean_bps, STDDEV(markout_bps) AS std_bps FROM ( SELECT aggressor, (future_mid - mid_at_trade) / mid_at_trade * 10000 AS markout_bps FROM trade_mids tm LEFT JOIN LATERAL ( SELECT mid_price AS future_mid FROM l2_snapshots WHERE coin = tm.coin AND exchange_ts_ms >= tm.exchange_ts_ms + 100 ORDER BY exchange_ts_ms LIMIT 1 ) ON TRUE WHERE mid_at_trade > 0 ) GROUP BY aggressor """ result = self.conn.execute(query, [coin.upper(), start_ts, end_ts]).fetchall() return { row[0]: {"count": row[1], "mean_bps": row[2], "std_bps": row[3]} for row in result } def stats(self) -> dict: """Get current database statistics.""" return { "l2_snapshots": self.conn.execute("SELECT COUNT(*) FROM l2_snapshots").fetchone()[0], "trades": self.conn.execute("SELECT COUNT(*) FROM trades").fetchone()[0], "funding": self.conn.execute("SELECT COUNT(*) FROM funding").fetchone()[0], "coins": [r[0] for r in self.conn.execute( "SELECT DISTINCT coin FROM l2_snapshots" ).fetchall()], "date_range": self.conn.execute( "SELECT MIN(exchange_ts_ms), MAX(exchange_ts_ms) FROM l2_snapshots" ).fetchone(), } def close(self): if self._conn: self._conn.close() self._conn = None def main(): p = argparse.ArgumentParser(description="DuckDB Tick Data Loader") p.add_argument("--data-dir", default="data/raw") p.add_argument("--db", default="data/normalized/ftdt_tick.db") p.add_argument("--coins", nargs="+", default=["BTC", "ETH"]) p.add_argument("--start-date", default="2026-01-01") p.add_argument("--end-date", default=date.today().isoformat()) p.add_argument("--stats", action="store_true", help="Print DB stats and exit") args = p.parse_args() logging.basicConfig(level=logging.INFO, format="%(asctime)s %(message)s", datefmt="%H:%M:%S") loader = DuckDBLoader(args.db) if args.stats: s = loader.stats() print(f"DuckDB: {args.db}") print(f" L2 snapshots: {s['l2_snapshots']:,}") print(f" Trades: {s['trades']:,}") print(f" Funding: {s['funding']:,}") print(f" Coins: {s['coins']}") print(f" Date range: {s['date_range']}") loader.close() return totals = loader.load_all(args.data_dir, args.coins, args.start_date, args.end_date) print(f"Loaded: {totals['l2']} L2 snapshots, {totals['trades']} trades, " f"{totals['funding']} funding observations") loader.close() if __name__ == "__main__": main()