fix: collector WebSocket trade parsing + DuckDB 1.5 API compatibility
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
This commit is contained in:
@@ -268,12 +268,12 @@ class HyperliquidCollector:
|
|||||||
)
|
)
|
||||||
self._latency.record_transport(time_ms, local_ts)
|
self._latency.record_transport(time_ms, local_ts)
|
||||||
|
|
||||||
async def _handle_trades(self, data: dict, local_ts: float):
|
async def _handle_trades(self, data: dict | list, local_ts: float):
|
||||||
coin = data.get("coin", "?")
|
|
||||||
if coin not in self._coins:
|
|
||||||
return
|
|
||||||
trade_list = data if isinstance(data, list) else [data]
|
trade_list = data if isinstance(data, list) else [data]
|
||||||
for trade in trade_list:
|
for trade in trade_list:
|
||||||
|
coin = trade.get("coin", "?") if isinstance(trade, dict) else "?"
|
||||||
|
if coin not in self._coins:
|
||||||
|
continue
|
||||||
time_ms = normalize_timestamp(trade.get("time"), source="hl")
|
time_ms = normalize_timestamp(trade.get("time"), source="hl")
|
||||||
self._store.push(
|
self._store.push(
|
||||||
channel="trades",
|
channel="trades",
|
||||||
@@ -311,7 +311,11 @@ class HyperliquidCollector:
|
|||||||
async with session.post(self._api_url, json={"type": "metaAndAssetCtxs"}, timeout=aiohttp.ClientTimeout(total=10)) as resp:
|
async with session.post(self._api_url, json={"type": "metaAndAssetCtxs"}, timeout=aiohttp.ClientTimeout(total=10)) as resp:
|
||||||
data = await resp.json()
|
data = await resp.json()
|
||||||
if isinstance(data, list) and len(data) >= 2:
|
if isinstance(data, list) and len(data) >= 2:
|
||||||
universe = data[0].get("universe", [])
|
raw_universe = data[0]
|
||||||
|
if isinstance(raw_universe, dict):
|
||||||
|
universe = raw_universe.get("universe", [])
|
||||||
|
else:
|
||||||
|
universe = raw_universe if isinstance(raw_universe, list) else []
|
||||||
ctxs = data[1]
|
ctxs = data[1]
|
||||||
now_ms = int(time.time() * 1000)
|
now_ms = int(time.time() * 1000)
|
||||||
for i, asset_info in enumerate(universe):
|
for i, asset_info in enumerate(universe):
|
||||||
@@ -357,7 +361,11 @@ class HyperliquidCollector:
|
|||||||
async with session.post(self._api_url, json={"type": "metaAndAssetCtxs"}, timeout=aiohttp.ClientTimeout(total=10)) as resp:
|
async with session.post(self._api_url, json={"type": "metaAndAssetCtxs"}, timeout=aiohttp.ClientTimeout(total=10)) as resp:
|
||||||
data = await resp.json()
|
data = await resp.json()
|
||||||
if isinstance(data, list) and len(data) >= 2:
|
if isinstance(data, list) and len(data) >= 2:
|
||||||
universe = data[0].get("universe", [])
|
raw_universe = data[0]
|
||||||
|
if isinstance(raw_universe, dict):
|
||||||
|
universe = raw_universe.get("universe", [])
|
||||||
|
else:
|
||||||
|
universe = raw_universe if isinstance(raw_universe, list) else []
|
||||||
ctxs = data[1]
|
ctxs = data[1]
|
||||||
now_ms = int(time.time() * 1000)
|
now_ms = int(time.time() * 1000)
|
||||||
for i, asset_info in enumerate(universe):
|
for i, asset_info in enumerate(universe):
|
||||||
|
|||||||
+9
-12
@@ -197,10 +197,9 @@ class DuckDBLoader:
|
|||||||
))
|
))
|
||||||
|
|
||||||
if rows:
|
if rows:
|
||||||
import duckdb
|
self.conn.executemany(
|
||||||
rel = duckdb.from_sequence(rows)
|
"INSERT OR IGNORE INTO l2_snapshots VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)",
|
||||||
self.conn.execute(
|
rows
|
||||||
"INSERT OR IGNORE INTO l2_snapshots SELECT * FROM rel"
|
|
||||||
)
|
)
|
||||||
logger.info("Loaded %d L2 snapshots for %s", len(rows), coin)
|
logger.info("Loaded %d L2 snapshots for %s", len(rows), coin)
|
||||||
|
|
||||||
@@ -242,10 +241,9 @@ class DuckDBLoader:
|
|||||||
))
|
))
|
||||||
|
|
||||||
if rows:
|
if rows:
|
||||||
import duckdb
|
self.conn.executemany(
|
||||||
rel = duckdb.from_sequence(rows)
|
"INSERT INTO trades VALUES (?,?,?,?,?,?,?)",
|
||||||
self.conn.execute(
|
rows
|
||||||
"INSERT INTO trades SELECT * FROM rel"
|
|
||||||
)
|
)
|
||||||
logger.info("Loaded %d trades for %s", len(rows), coin)
|
logger.info("Loaded %d trades for %s", len(rows), coin)
|
||||||
|
|
||||||
@@ -282,10 +280,9 @@ class DuckDBLoader:
|
|||||||
))
|
))
|
||||||
|
|
||||||
if rows:
|
if rows:
|
||||||
import duckdb
|
self.conn.executemany(
|
||||||
rel = duckdb.from_sequence(rows)
|
"INSERT OR IGNORE INTO funding VALUES (?,?,?,?,?,?)",
|
||||||
self.conn.execute(
|
rows
|
||||||
"INSERT OR IGNORE INTO funding SELECT * FROM rel"
|
|
||||||
)
|
)
|
||||||
logger.info("Loaded %d funding observations for %s", len(rows), coin)
|
logger.info("Loaded %d funding observations for %s", len(rows), coin)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user