第 9 章 · 02 BTC 地址 WebSocket 监控(btcagent) 本节摘要:本节精读 (约 242 行)——一个把「实时数据流」与「LLM 分析」缝在一起的实验脚本。它用 连上 ,订阅一个 BTC 地址,每当这个地址有交易发生(进账或出账),WebSocket 推送一条消息,脚本就把它喂给一个 swarms ,让 GPT-4o 现场分析「这笔交易有什么意义、有没有风险、网络流向如何」。本节重点讲三个工程主题:WebSocket 长连接的「订阅→推送」模型、为什么用独立线程 + signal.pause 跑 WebSocket(避免阻塞主线程)、以及信号处理(SIGINT/SIGTERM)如何实现优雅关闭。
本节摘要:本节精读
experimental/btc_agent.py(约 242 行)——一个把「实时数据流」与「LLM 分析」缝在一起的实验脚本。它用websocket-client连上wss://ws.blockchain.info/inv,订阅一个 BTC 地址,每当这个地址有交易发生(进账或出账),WebSocket 推送一条消息,脚本就把它喂给一个 swarmsAgent,让 GPT-4o 现场分析「这笔交易有什么意义、有没有风险、网络流向如何」。本节重点讲三个工程主题:WebSocket 长连接的「订阅→推送」模型、为什么用独立线程 + signal.pause 跑 WebSocket(避免阻塞主线程)、以及信号处理(SIGINT/SIGTERM)如何实现优雅关闭。这段代码展示了「事件驱动 + LLM」的组合范式——每条事件触发一次 Agent 推理,适合监控告警、异常检测类场景。
内容来源:原项目源码
experimental/btc_agent.py,逐行精读并套用体系化模板。
⚠️ 现实澄清:
wss://ws.blockchain.info/inv是 Blockchain.com 的公开实时交易流,免费但不稳定(常断连、限频)。代码里的重连逻辑很朴素。生产级监控会用付费的链上数据 API(如 QuickNode、Helius)。这段代码的价值在范式展示,不在稳定运行。
阅读完本节,你应当能够:
addr_sub 订阅、utx/tx 推送)。BTCTransactionMonitor 的四个回调(_on_open/_on_message/_on_error/_on_close)。signal.pause。监控「某个地址何时有新交易」这种需求,HTTP 不合适:
所以实时交易流(Blockchain.com、币安成交流、Solana 链上事件)都用 WebSocket。客户端只做两件事:连上 + 订阅,之后被动接收推送。
Blockchain.com 的 /inv 端点用一个简单的 JSON「op」协议:
{"op": "addr_sub", "addr": "<地址>"}:订阅某地址,之后该地址有任何相关交易都会被推过来。{"op": "utx" 或 "tx", "x": {...交易数据...}}:有新交易时推送,x 字段是交易详情(哈希、输入、输出、金额等)。本脚本在 _on_open(连接建立时)发订阅,在 _on_message(收到推送时)解析交易。
文件开头(第 13-52 行)建模型和 Agent:
model = OpenAIChat( model_name="gpt-4o", openai_api_key=os.getenv("OPENAI_API_KEY") ) BTC_AGENT_SYSTEM_PROMPT = """You are a specialized Bitcoin transaction analysis agent. Your role is to analyze new Bitcoin transactions in real-time and provide insights about: 1. Transaction significance and context 2. Pattern analysis and behavioral indicators 3. Risk assessment and unusual characteristics ... """ class BTCTransactionMonitor: def __init__(self): self.agent = Agent( agent_name="BTC-Analysis-Agent", system_prompt=BTC_AGENT_SYSTEM_PROMPT, llm=model, max_loops="auto", autosave=True, ... context_length=4000, ) self.running = False self.ws = None self.monitored_address = None
观察:
context_length=4000:上下文不大,够放一笔交易数据 + system prompt。BTC_AGENT_SYSTEM_PROMPT 让 Agent 从「交易意义/模式/风险/网络流向/经济影响」五个维度分析。running(是否在跑)、ws(连接对象)、monitored_address(当前订阅的地址)。websocket.WebSocketApp 用回调驱动,本脚本定义了四个(第 159-167 行组装):
def _connect_websocket(self): self.ws = websocket.WebSocketApp( "wss://ws.blockchain.info/inv", on_message=self._on_message, on_error=self._on_error, on_close=self._on_close, on_open=self._on_open, )
逐个看。
def _on_open(self, ws): logger.info("WebSocket connection established") subscription = { "op": "addr_sub", "addr": self.monitored_address, } ws.send(json.dumps(subscription))
连接一建立,立刻发 addr_sub 订阅目标地址。这一步错过就收不到推送。
def _on_message(self, ws, message): try: data = json.loads(message) if data.get("op") == "utx" or data.get("op") == "tx": tx_data = data.get("x", {}) # 收集这笔交易涉及的所有地址(输入+输出) addresses = [] for out in tx_data.get("out", []): addresses.append(out.get("addr", "")) for inp in tx_data.get("inputs", []): prev_out = inp.get("prev_out", {}) addresses.append(prev_out.get("addr", "")) # 命中我订阅的地址才分析 if self.monitored_address in addresses: logger.info(f"New transaction detected: {tx_data.get('hash', 'Unknown')}") analysis = self.analyze_transaction(tx_data) ... self._store_analysis(tx_data.get("hash", "unknown"), {...}) except json.JSONDecodeError: logger.error("Failed to decode websocket message") except Exception as e: logger.error(f"Error processing message: {str(e)}")
关键逻辑:
utx(unconfirmed tx)和 tx(confirmed)。self.analyze_transaction(tx_data) 把交易喂给 Agent。analysis_<hash>.json。def analyze_transaction(self, tx_data): value_btc = sum(out.get("value", 0) for out in tx_data.get("out", [])) / 100000000.0 analysis_prompt = f""" New Bitcoin transaction detected: Transaction Hash: {tx_data.get('hash', 'Unknown')} Time: {datetime.fromtimestamp(tx_data.get('time', 0))} Value: {value_btc} BTC Inputs: {len(tx_data.get('inputs', []))} Outputs: {len(tx_data.get('out', []))} ... """ return self.agent.run(analysis_prompt)
注意金额换算:
/ 100000000.0 转成 BTC。💡 事件驱动 + LLM 的范式:每来一条事件,把事件结构化成文本,喂给 Agent 让它「解读」。这种模式适合:链上监控(大额转账告警)、日志异常检测、IoT 数据流解读。优点是 LLM 能处理非结构化判断;缺点是成本与延迟——每条事件一次 LLM 调用,高频流会烧钱+延迟堆积。
def _on_error(self, ws, error): logger.error(f"WebSocket error: {str(error)}") def _on_close(self, ws, close_status_code, close_msg): logger.info("WebSocket connection closed") if self.running: logger.info("Attempting to reconnect...") self._connect_websocket() # ← 关闭后自动重连(但没重新 run_forever,有 bug)
注意 _on_close 的重连逻辑有缺陷:它重建了 WebSocketApp 对象,但没重新调 run_forever(见下一节线程模型),所以重连其实不会真的恢复事件循环。这是实验代码的典型瑕疵。
monitor_address(第 185-211 行)是启动入口:
def monitor_address(self, address: str): self.monitored_address = address self.running = True logger.info(f"Starting real-time monitoring for address: {address}") # 注册信号处理 signal.signal(signal.SIGINT, self._handle_shutdown) signal.signal(signal.SIGTERM, self._handle_shutdown) # WebSocket 跑在独立线程 self._connect_websocket() ws_thread = threading.Thread(target=self.ws.run_forever) ws_thread.daemon = True # ← daemon:主线程退出时自动结束 ws_thread.start() # 主线程保持存活 while self.running: signal.pause() # ← 主线程睡到收到信号
为什么这么写:
ws.run_forever() 是阻塞的:它内部跑事件循环,不返回。如果在主线程调,主线程就被占住,没法响应信号。run_forever 在子线程跑,主线程解放出来。signal.pause():把主线程挂起,直到收到信号(Ctrl+C / kill)。这是让 Python 进程「活着但不忙等」的常见技巧。⚠️ 平台限制:
signal.pause()是 Unix-only,Windows 上不存在。本脚本在 Windows 跑会抛AttributeError。要跨平台,Windows 应换成threading.Event().wait()。这是实验代码未考虑跨平台的又一例。
_handle_shutdown(第 213-217 行):
def _handle_shutdown(self, signum, frame): logger.info("Shutting down...") self.stop() sys.exit(0) def stop(self): self.running = False if self.ws: self.ws.close() logger.info("Monitoring stopped")
stop() 把 running 置 False、关闭 WS。sys.exit(0) 让进程干净退出。这套「signal → handler → stop → exit」是长跑服务的标准收尾模式,避免粗暴 kill 留下脏状态(如未关闭的连接、未刷盘的日志)。
这种「每事件一次 LLM」模式的致命问题是成本:
所以这套模式只适合低频地址:冷钱包、大额监控(一天几笔那种)。高频流要么换小模型(4o-mini)、要么先用规则过滤「只把可疑交易喂 LLM」、要么用流式聚合(攒一批再分析)。本脚本没做任何过滤,每笔都喂,生产不可行。
💡 生产化思路:加一层规则引擎——先用代码判断「这笔是否值得让 LLM 看」(如金额 > 阈值、地址在黑名单、模式异常),只有可疑的才喂 Agent。这样把 LLM 调用从「每笔」降到「每少数」,成本和延迟才可控。这是「LLM 作为稀有资源」的正确用法——别让它在每条事件上都跑。
可借鉴:
局限:
_on_close 重连没重新 run_forever。signal.pause() Unix-only。addr_sub 订阅,服务端推 utx/tx,带 x 交易详情。_on_open 发订阅;_on_message 解析+过滤+触发 Agent;_on_error/_on_close 处理异常(重连有 bug)。run_forever 阻塞,放 daemon 子线程;主线程 signal.pause() 挂起等信号。_handle_shutdown → stop() 关 WS → exit,优雅收尾。下一节,我们看
experimental/的最后一段——crypto_agent_wrapper.py,它用「适配器模式」封装一个外部 cryptoagent 库,展示了对外部依赖集成的思路,以及依赖缺失时的局限。