00bb06434e
data/collectors/hyperliquid.py: - Fixed _handle_trades to handle HL WebSocket trade data as list[dict] (each trade dict has its own 'coin' field) instead of assuming single dict - Fixed _poll_funding and _poll_open_interest to guard against data[0] being list vs dict (mainnet vs testnet response format difference) data/duckdb_load.py: - Replaced deprecated duckdb.from_sequence() with executemany() for l2_snapshots, trades, and funding table inserts (DuckDB 1.5.x API) Tick runner verified on real mainnet BTC data: - 324 events (45 L2 + 279 trades), VPIN-gated A-S MM, 1 trade, 98 cancels - HFT dashboard generated with 8 panels from real tick data
497 lines
16 KiB
Python
497 lines
16 KiB
Python
"""
|
|
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:
|
|
self.conn.executemany(
|
|
"INSERT OR IGNORE INTO l2_snapshots VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)",
|
|
rows
|
|
)
|
|
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:
|
|
self.conn.executemany(
|
|
"INSERT INTO trades VALUES (?,?,?,?,?,?,?)",
|
|
rows
|
|
)
|
|
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:
|
|
self.conn.executemany(
|
|
"INSERT OR IGNORE INTO funding VALUES (?,?,?,?,?,?)",
|
|
rows
|
|
)
|
|
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."""
|
|
self.init_schema()
|
|
l2_count = self.conn.execute(
|
|
"SELECT COUNT(*) FROM l2_snapshots"
|
|
).fetchone()[0]
|
|
trade_count = self.conn.execute(
|
|
"SELECT COUNT(*) FROM trades"
|
|
).fetchone()[0]
|
|
funding_count = self.conn.execute(
|
|
"SELECT COUNT(*) FROM funding"
|
|
).fetchone()[0]
|
|
coins_result = self.conn.execute(
|
|
"SELECT DISTINCT coin FROM l2_snapshots"
|
|
).fetchall()
|
|
date_result = self.conn.execute(
|
|
"SELECT MIN(exchange_ts_ms), MAX(exchange_ts_ms) FROM l2_snapshots"
|
|
).fetchone()
|
|
return {
|
|
"l2_snapshots": l2_count,
|
|
"trades": trade_count,
|
|
"funding": funding_count,
|
|
"coins": [r[0] for r in coins_result],
|
|
"date_range": date_result,
|
|
}
|
|
|
|
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()
|