"""修正 Hummingbot bitget_perpetual candles feed 的换根延迟。 上游 _parse_websocket_message 里是: candle = data["data"][0] 而 Bitget 在换根时会推一条带两根的消息 [上一根, 新一根]。取 [0] 拿到的是 上一根,其时间戳与 deque 尾部相同,于是只做了原地更新;新一根要等下一条 单元素消息才进入 deque——实测晚约 1.1 秒。 不能简单改成 [-1]:那样上一根的收盘价就永远停在换根前约 1 秒的那次推送上。 1m 信号对 0.25bp 的扰动都会换掉一半(见 signal_sensitivity.py),收盘价 偏一个 tick 是不能接受的。所以这里把**除最后一根外的元素就地写回 deque**, 再把最后一根交给基类走正常的 append 流程。 已向上游反馈前,本地用子类覆盖,不改动镜像。 ## 为什么必须有启动断言 子类覆盖的失效方式是**静默**的:上游若把 `_parse_websocket_message` 改名、 或改走别的钩子,我们的覆盖就成了死代码,行情悄悄退回慢 1.06 秒,不崩、 不报错、不留日志,只会让收益慢慢变差,几周后才从统计里看出来。 `assert_patch_effective()` 不做名字检查——名字对不上未必失效,名字对得上 也未必生效。它喂一条合成的两元素消息,直接验证行为:基类返回首元素(bug 仍在、覆盖仍有必要),子类返回末元素(覆盖确实生效)。再加一条源码检查 确认基类的收包循环还在调这个钩子。任一不满足就在启动时抛错。 """ from __future__ import annotations import inspect from typing import Any, Dict, Optional import numpy as np from hummingbot.data_feed.candles_feed.bitget_perpetual_candles import ( BitgetPerpetualCandles, ) from hummingbot.data_feed.candles_feed.candles_base import CandlesBase def _row_to_dict(row: list, ensure_s) -> Dict[str, Any]: return {"timestamp": ensure_s(int(row[0])), "open": float(row[1]), "high": float(row[2]), "low": float(row[3]), "close": float(row[4]), "volume": float(row[5]), "quote_asset_volume": float(row[6]), "n_trades": 0., "taker_buy_base_volume": 0., "taker_buy_quote_volume": 0.} class PatchedBitgetPerpetualCandles(BitgetPerpetualCandles): """与上游唯一的差别:一条消息里的多根 K 线全部处理,而非只取第一根。""" def _parse_websocket_message(self, data: dict) -> Optional[Dict[str, Any]]: if data == "pong": return None if not (data and data.get("data") and data.get("action") == "update"): return None rows = data["data"] # 前面的元素都是已收盘 K 线的最终值:就地覆盖,保住真实收盘价 for row in rows[:-1]: d = _row_to_dict(row, self.ensure_timestamp_in_seconds) self._overwrite_existing(d) # 最后一根交给基类:时间戳更大就 append,相同就原地更新 return _row_to_dict(rows[-1], self.ensure_timestamp_in_seconds) def _overwrite_existing(self, d: Dict[str, Any]) -> None: if not len(self._candles): return ts = int(d["timestamp"]) if int(self._candles[-1][0]) != ts: return self._candles[-1] = np.array( [d["timestamp"], d["open"], d["high"], d["low"], d["close"], d["volume"], d["quote_asset_volume"], d["n_trades"], d["taker_buy_base_volume"], d["taker_buy_quote_volume"]] ).astype(float) # 换根时 Bitget 推的就是这个形状:[上一根, 新一根] _PROBE = { "action": "update", "arg": {"instType": "USDT-FUTURES", "channel": "candle1m", "instId": "BTCUSDT"}, "data": [ ["1700000040000", "1", "1", "1", "1", "1", "1", "1"], ["1700000100000", "2", "2", "2", "2", "2", "2", "2"], ], } def assert_patch_effective() -> None: """启动即验证覆盖真的生效,否则抛错。让静默失效变成启动失败。""" src = inspect.getsource(CandlesBase._process_websocket_messages_task) if "_parse_websocket_message" not in src: raise RuntimeError( "上游收包循环已不再调用 _parse_websocket_message," "patched_candles 的覆盖失效。需重新定位钩子后再启动。") stock = BitgetPerpetualCandles("BTC-USDT", "1m", 20) ours = PatchedBitgetPerpetualCandles("BTC-USDT", "1m", 20) got_stock = stock._parse_websocket_message(_PROBE) got_ours = ours._parse_websocket_message(_PROBE) head_ts = stock.ensure_timestamp_in_seconds(int(_PROBE["data"][0][0])) tail_ts = stock.ensure_timestamp_in_seconds(int(_PROBE["data"][-1][0])) if not got_ours or int(got_ours["timestamp"]) != int(tail_ts): raise RuntimeError( f"覆盖未生效:子类返回 {got_ours and got_ours.get('timestamp')}," f"应为末元素 {tail_ts}。") if got_stock and int(got_stock["timestamp"]) == int(tail_ts): # 上游自己修好了。此时覆盖无害但已多余,明确说出来,免得以后 # 有人以为那 1.06 秒还是靠我们拿回来的 print(" [补丁] 上游已自行修正换根解析,本地覆盖现为冗余,可移除", flush=True) elif not got_stock or int(got_stock["timestamp"]) != int(head_ts): raise RuntimeError( f"基类行为与预期不符:返回 " f"{got_stock and got_stock.get('timestamp')}," f"既非首元素 {head_ts} 也非末元素 {tail_ts}。" f"上游改了解析逻辑,补丁的前提需重新确认。") else: print(f" [补丁] 覆盖生效:基类取首元素 {int(head_ts)}、" f"本地取末元素 {int(tail_ts)}", flush=True)