diff --git a/data/collectors/hyperliquid.py b/data/collectors/hyperliquid.py index 5e7d0fe..9de0ba0 100644 --- a/data/collectors/hyperliquid.py +++ b/data/collectors/hyperliquid.py @@ -268,12 +268,12 @@ class HyperliquidCollector: ) self._latency.record_transport(time_ms, local_ts) - async def _handle_trades(self, data: dict, local_ts: float): - coin = data.get("coin", "?") - if coin not in self._coins: - return + async def _handle_trades(self, data: dict | list, local_ts: float): trade_list = data if isinstance(data, list) else [data] 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") self._store.push( 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: data = await resp.json() 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] now_ms = int(time.time() * 1000) 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: data = await resp.json() 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] now_ms = int(time.time() * 1000) for i, asset_info in enumerate(universe): diff --git a/data/duckdb_load.py b/data/duckdb_load.py index 2df4844..91b14c4 100644 --- a/data/duckdb_load.py +++ b/data/duckdb_load.py @@ -197,10 +197,9 @@ class DuckDBLoader: )) if rows: - import duckdb - rel = duckdb.from_sequence(rows) - self.conn.execute( - "INSERT OR IGNORE INTO l2_snapshots SELECT * FROM rel" + self.conn.executemany( + "INSERT OR IGNORE INTO l2_snapshots VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)", + rows ) logger.info("Loaded %d L2 snapshots for %s", len(rows), coin) @@ -242,10 +241,9 @@ class DuckDBLoader: )) if rows: - import duckdb - rel = duckdb.from_sequence(rows) - self.conn.execute( - "INSERT INTO trades SELECT * FROM rel" + self.conn.executemany( + "INSERT INTO trades VALUES (?,?,?,?,?,?,?)", + rows ) logger.info("Loaded %d trades for %s", len(rows), coin) @@ -282,10 +280,9 @@ class DuckDBLoader: )) if rows: - import duckdb - rel = duckdb.from_sequence(rows) - self.conn.execute( - "INSERT OR IGNORE INTO funding SELECT * FROM rel" + self.conn.executemany( + "INSERT OR IGNORE INTO funding VALUES (?,?,?,?,?,?)", + rows ) logger.info("Loaded %d funding observations for %s", len(rows), coin)