入账监听实现指南
为什么需要扫块
「监听某个地址有没有收到钱」看起来简单,实际上有三种做法,只有一种是对的。
- 轮询余额(错):定时查地址余额,发现变多就算入账。这样你拿不到 txid、拿不到付款方、分不清是两笔还是一笔,也无法对账。金额相同的两笔转账会被合并成一笔,客户投诉时你无从查起。
- 依赖第三方索引接口(受制于人):用交易列表类接口按地址查历史。这类接口是索引服务而非节点原生能力,有限流、有延迟、有停服风险,且你无法验证它给的数据是否完整——漏推一笔你不会知道。
- 扫块(正确):顺序读取每一个区块,自己从中筛出与你相关的转账。区块是链上的完整事实,只要不跳块,就不可能漏单;每一笔都带 txid、块高、付款方,天然可对账、可复核、可重放。
本文给出一套可直接落地的扫块实现:同时支持原生 TRX 与 TRC20(USDT),包含假充值防御、二次核验和完整状态机,并且在设计上把节点请求量压到极低——稳态下每分钟只需 0.2 次请求,远低于任何限流阈值。
https://api.trxapi.io?apikey=YOUR_API_KEY1M / day · 5 / s关键常量速查
实现前先把这几个值抄下来。它们写错都不会报错——扫描器照常运行、日志干净,只是永远扫不到任何入账,或者把不该入账的当成入账。
| 常量 | 值 | 用途与注意 |
|---|---|---|
| USDT 合约地址(区块中) | 41a614f803b6fd780986a42c78ec9c7f77e6ded13c | 与 contract_address 全等比对。带 41 前缀 |
| USDT 合约地址(日志中) | a614f803b6fd780986a42c78ec9c7f77e6ded13c | 与 log.address 全等比对。不带 41 前缀,与上一行是两个常量 |
| USDT Base58 地址 | TR7NHqjeKQxGTCi8q8ZY4pL8otSzgjLj6t | 仅供人工核对与展示;代码里不要用它做比对 |
| transfer 方法选择器 | a9059cbb | calldata 的前 8 位十六进制 |
| Transfer 事件签名 | ddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef | 与 topics[0] 全等比对,且 topics 必须正好 3 个 |
| USDT 小数位 | 6 | 写死在配置里,不要运行时调 decimals() |
| 单批拉块上限 | 100 | 超过会静默返回空数组,不报错 |
| 入账确认数 | 20 个区块 | solidity 端点实测落后 fullnode 18 个块,留 2 块余量 |
41 前缀、在日志里不带,必须准备两个常量,混用的症状是「一笔入账都扫不到」,看起来像没人转账。确认数不要自己数——以 solidity 端点查得到为准,那是超级节点已确认、不可回滚的数据;如果你的风控流程必须落一个数字,不得低于 20,低于 18 等于根本没等固化。换测试网或换代币时,这张表里的每一个值都要跟着换,尤其是合约地址:拿主网 USDT 地址去测试网扫,你会一笔都扫不到,且没有任何报错。架构总览
三个组件,各司其职。把职责搞混是这类系统最常见的事故来源,尤其是把状态放进 Redis。
- 节点 API——唯一的链上事实来源。批量拉块、拉回执、做二次核验。
- Redis——只放可以随时丢弃并重建的东西:监听地址集合、去重前置缓存、待推送队列、区块缓存。它是加速器,不是账本。
- 数据库(PostgreSQL / MySQL)——唯一的账。扫描水位、入账记录、状态机全部落库,靠唯一索引保证幂等。
数据流:
# ┌── getblockbylimitnext ── 一次 100 个块 ────────────┐ # │ ▼ # │ 本地解析:TransferContract / TriggerSmartContract # │ │ # │ ▼ Redis SET 过滤(零请求) # │ 命中监听地址? # │ 否 ──► 丢弃(99.99% 走这条) # │ 是 ──► gettransactioninfobyblocknum # │ │ 确认 receipt + Transfer 日志 # │ ▼ # │ 落库 status=detected(唯一索引幂等) # │ │ # │ ▼ 等 20 个块固化 # │ solidity 二次核验 ──► status=confirmed ──► 入账 # └── 水位 +100 落库,继续下一批 ─────────────────────┘
配额与批量拉块
这是整套设计里最关键、也最容易做错的一环。逐块请求会直接把你的配额打爆,而正确做法能把请求量降到它的两百分之一。
TRON 每 3 秒出一个块,每分钟 20 个。如果按「每块请求一次区块 + 一次回执」来写:
| 做法 | 节点请求量 | 免密公共 RPC | 注册客户(apikey) |
|---|---|---|---|
| 逐块拉取(区块 + 回执) | 40 / min | ✗ 超限 | ✓ 可用 |
| 批量拉块 + 按需拉回执 | 0.2 / min | ✓ 可用 | ✓ 余量充足 |
getblockbylimitnext 允许一次请求取回最多 100 个区块。100 个块 = 5 分钟的链上数据,也就是说保持实时同步平均每 5 分钟才需要发一次请求。而回执(gettransactioninfobyblocknum)只在这一批块里确实命中了你的监听地址时才去拉——绝大多数批次一笔都不命中,一次回执请求都不用发。
按注册客户 5 次/秒、每天 100 万次的免费额度算:实时同步一天只消耗约 300~600 次,占额度的 0.05%。剩下的额度足够你同时跑历史回补:全速回补时每秒可处理约 250 个块,等于每秒消化 12 分钟的链上数据,补一整年的历史约需 20 小时。
getblockbylimitnext 的区间是左闭右开,且超过 100 个块时不会报错,而是静默返回空数组。写成 {"startNum":n,"endNum":n+200} 你会得到 {}——没有异常、没有错误码,扫描器会安静地什么都不做,看起来像「已经追平了」。请务必把每批限制在 100 以内,并在代码里断言返回块数与请求区间一致。扫块与水位
永远不要写「定时取最新块」。取最新块的写法在节点抖动、进程 GC、网络超时的瞬间就会跳过区块,而且不会抛任何错误——你只会在客户投诉「转了钱没到账」时才发现,且无法知道漏了哪些。
正确形态是水位追块:数据库里存「已处理到第几块」,每轮从水位往前推进,追平了才休息。节点慢了你自然落后、下轮补上;节点快了你连续追。任何情况下都不会跳块。
水位必须落在数据库里,不能只存 Redis——Redis 一次 flush 或一次内存淘汰,你要么全链重扫,要么永久漏掉一段。
# Correct: watermark-driven catch-up. Never "fetch the latest block". BATCH = 100 # hard cap of getblockbylimitnext cursor = db_load_cursor() # last processed block, from DB while True: head = rpc("wallet", "getnowblock")["block_header"]["raw_data"]["number"] if cursor >= head: time.sleep(1.5) # caught up — idle, costs almost nothing continue start = cursor + 1 end = min(start + BATCH, head + 1) # endNum is EXCLUSIVE blocks = rpc("wallet", "getblockbylimitnext", {"startNum": start, "endNum": end}).get("block", []) if len(blocks) != end - start: # never trust a short read log.warning("short read: want %d got %d", end - start, len(blocks)) time.sleep(2); continue # retry, do NOT advance the cursor process(blocks) cursor = end - 1 db_save_cursor(cursor) # advance only after the batch is durable
原生 TRX 入账
原生 TRX 转账在区块里是 TransferContract,所有需要的信息都直接在区块数据里,不需要额外请求:
parameter.value.to_address—— 收款地址,hex 格式且带41前缀parameter.value.owner_address—— 付款地址parameter.value.amount—— 金额,单位 SUN(1 TRX = 1,000,000 SUN)ret[0].contractRet—— 执行结果,必须等于 SUCCESS
两个容易踩的点:
- 地址是 hex 不是 Base58。区块里是
41dc6b4b…这种形式,要转成T开头的 Base58 才能和你数据库里的地址比对。转换是41 前缀 + 20 字节地址再做 Base58Check(双 SHA256 取前 4 字节做校验位)。 - 合约内部转账不是 TransferContract。如果付款方是通过合约给你打的 TRX,区块里看到的是
TriggerSmartContract,真正的转账藏在回执的internal_transactions里。只认 TransferContract 会漏掉这类到账。交易所场景建议一并处理。
# Native TRX — everything comes straight from the block, zero extra requests for tx in block.get("transactions", []): if tx.get("ret", [{}])[0].get("contractRet") != "SUCCESS": continue # failed on-chain — ignore c = tx["raw_data"]["contract"][0] if c["type"] != "TransferContract": continue v = c["parameter"]["value"] to = to_base58(v["to_address"]) # 41dc6b… -> T… if not redis.sismember("watch:addrs", to): continue # not ours — 99.99% land here hit(kind="TRX", txid=tx["txID"], to=to, sender=to_base58(v["owner_address"]), amount=int(v["amount"]), decimals=6) # amount is in SUN
TRC20(USDT)入账
TRC20 转账不是一种独立的交易类型,它是一次普通的合约调用(TriggerSmartContract)。区块里只能看到「谁调用了哪个合约、参数是什么」,看不到转账是否真的成功。
所以分两步走,这也正是把请求量压下来的关键:
- 从区块里筛候选(零额外请求)。判断三件事:合约地址是不是真 USDT、方法选择器是不是
a9059cbb(即transfer(address,uint256))、ret[0].contractRet是不是 SUCCESS。data字段的布局是固定的:8 位选择器 + 64 位收款地址(右对齐,取后 40 位)+ 64 位金额。 - 命中后才拉回执核对(按需请求)。用
gettransactioninfobyblocknum一次取回整块的执行信息,检查receipt.result并以 Transfer 日志里的金额为准——不要相信 calldata。calldata 是「请求转多少」,日志是「实际转了多少」,遇到手续费型代币或代理合约时两者并不相等。
USDT = "41a614f803b6fd780986a42c78ec9c7f77e6ded13c" # WITH 41 prefix (block) USDT_LOG = USDT[2:] # WITHOUT it (logs!) SELECTOR = "a9059cbb" # transfer(address,uint256) TRANSFER = "ddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef" # --- step 1: candidates from the block, no extra request --------------- c = tx["raw_data"]["contract"][0] v = c["parameter"]["value"] if c["type"] == "TriggerSmartContract" \ and v.get("contract_address") == USDT \ and v.get("data", "")[:8] == SELECTOR \ and tx.get("ret", [{}])[0].get("contractRet") == "SUCCESS": data = v["data"] to = to_base58("41" + data[8:72][-40:]) # arg is 20 bytes, add 41 if redis.sismember("watch:addrs", to): candidates.append((tx["txID"], to, int(data[72:136], 16))) # --- step 2: only for blocks with candidates, pull receipts ------------ infos = rpc("wallet", "gettransactioninfobyblocknum", {"num": block_num}) by_id = {i["id"]: i for i in infos} for txid, to, calldata_amt in candidates: info = by_id.get(txid) if not info or info.get("receipt", {}).get("result") != "SUCCESS": continue # fake deposit — drop it for lg in info.get("log", []): if lg.get("address") != USDT_LOG: # must be the REAL contract continue tp = lg.get("topics", []) if len(tp) != 3 or tp[0] != TRANSFER: # must be a standard Transfer continue if to_base58("41" + tp[2][-40:]) != to: continue amount = int(lg["data"], 16) # trust the LOG, not calldata hit(kind="USDT", txid=txid, to=to, amount=amount, decimals=6) break
41 前缀,日志里的 log.address 不带。区块中是 41a614f803…,日志中是 a614f803…。拿同一个常量去比对两处,其中一处必然永远不匹配——而症状是「一笔入账都扫不到」,看起来像是没人转账,非常难排查。合约判别与假币
任何人都可以在 TRON 上部署一个合约,把它叫做「USDT」、符号写成「USDT」、小数位设成 6,然后向你的地址转账 100 万枚。它会产生一条格式完全标准的 Transfer 日志。如果你按代币名称或符号识别,你就会给他入账 100 万 USDT。
唯一可靠的判据是合约地址——它由部署交易决定,无法伪造、无法抢注:
| 判据 | 正确做法 | 原因 |
|---|---|---|
| 合约地址 | 41a614f8…ded13cTR7NHqjeKQxGTCi8q8ZY4pL8otSzgjLj6t | 硬编码在配置里,全等比对。这是唯一不可伪造的身份 |
| 代币名称 / 符号 | 绝不使用 | 任何人可自定义,伪造成本为零 |
| 事件签名 | topics[0] == ddf252ad… | 确认这确实是 Transfer 事件,而非同名的自定义事件 |
| topics 数量 | len(topics) == 3 | 标准 Transfer 的 from / to 均为 indexed;数量不符说明是非标准实现,不应入账 |
| 小数位 | 按合约固定写死 | 不要运行时查询 decimals(),恶意合约可以对不同调用方返回不同值 |
二次核验与假充值
「假充值」是这个领域最经典的攻击,手法极其简单:攻击者发起一笔 USDT 转账,但故意不给足能量,交易会上链、会出现在区块里、会有 txid,但执行结果是 OUT_OF_ENERGY,代币一枚都没有动。如果你的系统只看「这笔交易在链上存在」就入账,他就凭空拿到了余额。
这不是理论风险。我们在写这篇文档时抽样了 TRON 主网连续 20 个区块:
也就是说每 100 笔 USDT 转账里就有大约 1 笔是失败的。这个比例足够高,任何漏掉这一步的系统上线当天就会被打穿。
防御是三层,缺一不可:
- 区块层:
ret[0].contractRet == "SUCCESS" - 回执层:
receipt.result == "SUCCESS",且金额以 Transfer 日志为准 - 固化层(真正入账前的最后一道):改用
walletsolidity端点重新核验一遍
第三层就是「匹配到之后再问一次节点」。它防的是另一类风险——链上重组。fullnode 返回的是最新数据,理论上可能被回滚;solidity 端点只返回已被超级节点确认、不可逆的数据。实测 solidity 端点落后 fullnode 约 18 个区块(≈54 秒),这就是你在「快」和「绝对安全」之间要付的代价。
核验要比对的不只是「查得到」,还包括块高是否一致、金额是否一致——这三项全对,才说明你扫到的那笔交易确实以你记录的形式被写进了不可逆的链上历史。
# Final check before crediting. Re-ask the node via the SOLIDITY endpoint: # solidified data cannot be rolled back. Runs once per deposit, ~20 blocks later. def verify_final(dep): if dep["kind"] == "TRX": tx = rpc("walletsolidity", "gettransactionbyid", {"value": dep["txid"]}) if not tx.get("txID"): return False, "not solidified yet, or reorged out" if tx["ret"][0]["contractRet"] != "SUCCESS": return False, "execution failed" v = tx["raw_data"]["contract"][0]["parameter"]["value"] if to_base58(v["to_address"]) != dep["to"] or int(v["amount"]) != dep["amount"]: return False, "recipient or amount mismatch" return True, "" info = rpc("walletsolidity", "gettransactioninfobyid", {"value": dep["txid"]}) if not info.get("id"): return False, "not solidified yet, or reorged out" if info.get("blockNumber") != dep["block"]: return False, "block height mismatch — possible reorg" if info.get("receipt", {}).get("result") != "SUCCESS": return False, "execution failed — fake deposit" total = sum(int(l["data"], 16) for l in info.get("log", []) if l.get("address") == USDT_LOG and len(l.get("topics", [])) == 3 and l["topics"][0] == TRANSFER and to_base58("41" + l["topics"][2][-40:]) == dep["to"]) if total != dep["amount"]: return False, "amount mismatch" return True, ""
入账状态机
入账不是一个布尔值,而是一条有向的状态链。只能向前推进,任何一步都不能跳过。把它写成状态机而不是一堆布尔字段,事故率会低一个数量级。
| 状态 | 含义 | 如何进入下一步 |
|---|---|---|
detected | 扫块命中,回执已确认执行成功 | 等待 20 个块固化。可以给用户展示「已检测到,确认中」,但绝不能动余额 |
confirmed | solidity 二次核验通过,不可逆 | 可以入账。这是唯一允许触发资金动作的状态 |
credited | 已写入用户余额 | 终态。入账与改状态必须在同一个数据库事务里 |
rejected | 核验不通过(执行失败 / 金额不符 / 重组) | 终态。金额不符要告警转人工,不要静默丢弃 |
orphaned | 超过阈值仍未固化 | 终态。记录并告警,供人工复核 |
两条铁律:
- 幂等靠数据库唯一索引,不靠代码判断。唯一键用
(txid, log_index)——一笔交易里可能包含多笔转账,只用 txid 会把它们合并成一笔。代码里的「先查再插」在并发下必然失效,唯一索引不会。 - 入账动作和状态变更必须原子。「加了余额但状态没改成 credited」意味着下次扫描会再加一次。放进同一个事务,用行锁保护余额。
detected 就给用户加钱——必须等 confirmed,展示可以早、动钱不行;② 用 is_confirmed / is_credited 两个布尔字段代替状态列,并发下两个都为真,钱加两次;③ rejected 静默丢弃——金额不符属于安全事件,必须告警到人。另外状态变更一律用条件更新(UPDATE ... SET status='credited' WHERE status='confirmed'),按受影响行数判断自己有没有抢到;「先查询再更新」在两个 worker 之间必然重复入账。Redis 配置与键设计
Redis 在这套系统里只做三件事:过滤加速、去重前置、推送队列。它存的每一样东西都必须能从数据库和链上重建。
配置里最要命的是淘汰策略:
# redis.conf — deposit monitor maxmemory 4gb maxmemory-policy noeviction # NEVER allkeys-lru / allkeys-random appendonly yes # the pending-notify queue must survive a restart appendfsync everysec bind 127.0.0.1 # never expose Redis to the internet requirepass <strong-password> rename-command FLUSHALL "" # flushing the watchlist = missing deposits
allkeys-lru 会在内存吃紧时把你的监听地址集合和去重标记一起淘汰掉。后果是漏单和重复入账,而且发生在半夜流量高峰、没有任何报错。用 noeviction:内存不够就报错,你能立刻发现并处理。如果确实要用 LRU,就用 volatile-lru,并且严格保证只有缓存类的键才设 TTL。键设计(前缀带版本号,方便整体重建):
# --- state-ish: no TTL, rebuildable from the database ----------------- v1:watch:addrs SET # every watched address; the hot path filter v1:notify:queue LIST # deposits pending push to your business system # --- cache: TTL is fine, losing it costs nothing ---------------------- v1:seen:{txid}:{idx} STRING TTL 3h # dedupe shortcut before hitting the DB v1:deposit:{address} ZSET TTL 3h # score=block, recent deposits per address v1:block:{num} STRING TTL 3h # raw block cache, for replay/debugging # memory: 100k watched addresses + 3h of hits ~= tens of MB. # caching raw blocks too: 3600 blocks * a few hundred KB ~= 1-2 GB.
v1:seen:* 只是省一次数据库查询的前置缓存,不是幂等判据。它过期了、被清了、Redis 整个挂了,都不能导致重复入账——真正的幂等永远是数据库那条唯一索引。设计时问自己一句:「Redis 现在整个清空,我的系统会不会多给用户一分钱?」答案必须是不会。数据表结构
两张表就够。金额用整数存最小单位(TRX 存 SUN、USDT 存 6 位小数的整数),绝对不要用浮点——浮点在金额上迟早会给你凑出一分钱的差额,而对账差一分钱和差一万块一样要查。
-- scan watermark: the single source of truth for "how far have I got" CREATE TABLE scan_cursor ( chain TEXT PRIMARY KEY, -- 'tron' block_num BIGINT NOT NULL, updated_at TIMESTAMPTZ NOT NULL DEFAULT now() ); CREATE TABLE deposits ( id BIGSERIAL PRIMARY KEY, txid TEXT NOT NULL, log_index INT NOT NULL DEFAULT 0, -- one tx can hold many transfers block_num BIGINT NOT NULL, block_time TIMESTAMPTZ NOT NULL, currency TEXT NOT NULL, -- 'TRX' | 'USDT' contract TEXT NOT NULL DEFAULT '', -- empty for native TRX from_addr TEXT NOT NULL, to_addr TEXT NOT NULL, amount NUMERIC(38,0) NOT NULL, -- smallest unit, integer only decimals INT NOT NULL, status TEXT NOT NULL, -- detected|confirmed|credited|rejected|orphaned reject_note TEXT NOT NULL DEFAULT '', created_at TIMESTAMPTZ NOT NULL DEFAULT now() ); -- THE idempotency guarantee. Not your application code. CREATE UNIQUE INDEX uq_deposits ON deposits (txid, log_index); CREATE INDEX ix_deposits_pending ON deposits (status, block_num); CREATE INDEX ix_deposits_addr ON deposits (to_addr, created_at DESC);
完整示例代码
下面是一个可直接运行的最小实现,只依赖标准库(Redis / 数据库部分留了明确的接入点)。把 API_KEY 换成你自己的即可跑起来。
#!/usr/bin/env python3 # TRON deposit monitor — native TRX + TRC20(USDT) # docs: https://www.trxapi.io/docs/deposit-monitor.html import hashlib, json, logging, time, urllib.request API_KEY = "YOUR_API_KEY" BASE = "https://api.trxapi.io/api/v1/openapi/node" USDT = "41a614f803b6fd780986a42c78ec9c7f77e6ded13c" USDT_LOG = USDT[2:] SELECTOR = "a9059cbb" TRANSFER = "ddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef" BATCH = 100 # hard cap of getblockbylimitnext CONFIRMS = 20 # solidity lags ~18 blocks; 20 gives margin log = logging.getLogger("deposit") def rpc(kind, method, body=None, retries=5): url = f"{BASE}/{kind}/{method}?apikey={API_KEY}" for i in range(retries): try: req = urllib.request.Request( url, method="POST", data=json.dumps(body or {}).encode(), headers={"Content-Type": "application/json"}) with urllib.request.urlopen(req, timeout=30) as r: return json.load(r) except Exception as e: wait = 2 ** i # 1s, 2s, 4s... on 429/timeout log.warning("rpc %s failed (%s), retry in %ds", method, e, wait) time.sleep(wait) raise RuntimeError(f"rpc {method} failed after {retries} tries") # ---- hex(41-prefixed) -> Base58Check (T...) -------------------------- B58 = "123456789ABCDEFGHJKLMNPQRSTUVWXYZabcdefghijkmnopqrstuvwxyz" def to_base58(hex41): raw = bytes.fromhex(hex41) chk = hashlib.sha256(hashlib.sha256(raw).digest()).digest()[:4] n = int.from_bytes(raw + chk, "big") out = "" while n: n, r = divmod(n, 58) out = B58[r] + out return out # ---- plug in your own storage --------------------------------------- def is_watched(addr): ... # redis.sismember("v1:watch:addrs", addr) def save_deposits(rows): ... # INSERT ... ON CONFLICT (txid,log_index) DO NOTHING def load_cursor(): ... # SELECT block_num FROM scan_cursor def save_cursor(n): ... # UPDATE scan_cursor SET block_num=$1 def scan_batch(start, end): """Fetch [start, end) in ONE request and extract everything we care about.""" blocks = rpc("wallet", "getblockbylimitnext", {"startNum": start, "endNum": end}).get("block", []) if len(blocks) != end - start: raise RuntimeError(f"short read: want {end-start}, got {len(blocks)}") hits = [] for blk in blocks: raw = blk["block_header"]["raw_data"] num, ts = raw.get("number", 0), raw["timestamp"] cands = [] for tx in blk.get("transactions", []): if tx.get("ret", [{}])[0].get("contractRet") != "SUCCESS": continue c = tx["raw_data"]["contract"][0] v = c["parameter"]["value"] if c["type"] == "TransferContract": # native TRX to = to_base58(v["to_address"]) if is_watched(to): hits.append(dict(kind="TRX", currency="TRX", contract="", txid=tx["txID"], log_index=0, block=num, ts=ts, to=to, sender=to_base58(v["owner_address"]), amount=int(v["amount"]), decimals=6)) elif c["type"] == "TriggerSmartContract": # TRC20 candidate data = v.get("data", "") if v.get("contract_address") != USDT: continue if data[:8] != SELECTOR or len(data) < 136: continue to = to_base58("41" + data[8:72][-40:]) if is_watched(to): cands.append((tx["txID"], to, to_base58(v["owner_address"]))) if cands: # receipts ONLY when needed hits += confirm_trc20(num, ts, cands) return hits def confirm_trc20(num, ts, cands): infos = rpc("wallet", "gettransactioninfobyblocknum", {"num": num}) by_id = {i["id"]: i for i in infos} out = [] for txid, to, sender in cands: info = by_id.get(txid) if not info or info.get("receipt", {}).get("result") != "SUCCESS": log.info("drop fake deposit %s (execution failed)", txid) continue for idx, lg in enumerate(info.get("log", [])): if lg.get("address") != USDT_LOG: # the real contract only continue tp = lg.get("topics", []) if len(tp) != 3 or tp[0] != TRANSFER: # a standard Transfer only continue if to_base58("41" + tp[2][-40:]) != to: continue out.append(dict(kind="USDT", currency="USDT", contract=USDT, txid=txid, log_index=idx, block=num, ts=ts, to=to, sender=sender, amount=int(lg["data"], 16), decimals=6)) return out def run(): cursor = load_cursor() while True: head = rpc("wallet", "getnowblock")["block_header"]["raw_data"]["number"] if cursor >= head: time.sleep(1.5) continue start = cursor + 1 end = min(start + BATCH, head + 1) try: hits = scan_batch(start, end) except RuntimeError as e: log.warning("batch %d-%d failed: %s", start, end, e) time.sleep(2) continue # cursor NOT advanced # save_deposits() and save_cursor() belong in ONE transaction save_deposits(hits) save_cursor(end - 1) cursor = end - 1 if hits: log.info("blocks %d-%d: %d deposits", start, end - 1, len(hits)) if __name__ == "__main__": logging.basicConfig(level=logging.INFO) run()
Node.js 版本的核心循环(其余部分同理):
// TRON deposit monitor — core loop (Node.js 18+, no dependencies) const API = 'https://api.trxapi.io/api/v1/openapi/node'; const KEY = 'YOUR_API_KEY'; const USDT = '41a614f803b6fd780986a42c78ec9c7f77e6ded13c'; const USDT_LOG = USDT.slice(2); const TRANSFER = 'ddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef'; async function rpc(kind, method, body = {}) { const r = await fetch(`${API}/${kind}/${method}?apikey=${KEY}`, { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify(body), }); if (r.status === 429) throw new Error('rate limited'); // back off return r.json(); } async function scanBatch(start, end) { // end is EXCLUSIVE const { block = [] } = await rpc('wallet', 'getblockbylimitnext', { startNum: start, endNum: end }); if (block.length !== end - start) throw new Error(`short read: want ${end - start}, got ${block.length}`); const hits = []; for (const blk of block) { const num = blk.block_header.raw_data.number ?? 0; const cands = []; for (const tx of blk.transactions ?? []) { if (tx.ret?.[0]?.contractRet !== 'SUCCESS') continue; const c = tx.raw_data.contract[0], v = c.parameter.value; if (c.type === 'TransferContract') { const to = toBase58(v.to_address); if (await isWatched(to)) hits.push({ kind: 'TRX', txid: tx.txID, logIndex: 0, block: num, to, sender: toBase58(v.owner_address), amount: BigInt(v.amount) }); } else if (c.type === 'TriggerSmartContract' && v.contract_address === USDT && (v.data ?? '').startsWith('a9059cbb')) { const to = toBase58('41' + v.data.slice(8, 72).slice(-40)); if (await isWatched(to)) cands.push({ txid: tx.txID, to }); } } if (cands.length) hits.push(...await confirmTrc20(num, cands)); } return hits; }
常见坑
- endNum 是开区间。想扫
[n, n+99]要传endNum = n+100。写错就是每批默默少扫一个块,一天漏掉 288 个块。 - 请求超过 100 个块返回空数组而不是报错。扫描器会安静地空转,监控上看一切正常。
- 区块里的合约地址带
41,日志里的log.address不带。混用会导致永远零命中。 - topics 里的地址是 32 字节右对齐的。要取后 40 位十六进制再补
41前缀,直接拿整个 topic 去转换会得到垃圾地址。 - 一笔交易可以包含多笔转账。幂等键必须是
(txid, log_index);只用 txid 会把批量转账吞掉,只入账第一笔。 - 金额别用浮点。USDT 的 6 位小数用 float 处理,迟早凑出 0.000001 的差额,而对账查这一分钱要花一整天。
- 不要在 detected 状态就入账。必须等 solidity 核验通过。展示可以早,动钱不行。
- 回补要限速。全速回补会瞬间打满配额并触发 429,也会拖垮你自己的数据库写入。给回补单独一个低速通道,并且可暂停。
- 别忘了合约内部转账。通过合约转给你的 TRX 不是 TransferContract,藏在回执的
internal_transactions里。 - Redis 清空不能影响正确性。如果清空 Redis 会导致重复入账或漏单,说明你把状态放错地方了。