系列导航

顺序文章内容
1Web3 链上转账完整流程(小白版)零代码,建立整体认识
2智能合约:从零看懂一个代币合约合约源码、ABI、测试、部署、转账时 EVM 里发生了什么
3前端:用 MetaMask + ethers 完成链上转账RPC、钱包、Provider/Signer、new Contract、转账 12 步
4📍当前:后端:监听链上事件并同步到数据库JSON-RPC 与 WebSocket、事件订阅、确认数与重组、幂等、可靠性

读完本篇你将学会

  • 说清楚 Web3 应用为什么还需要一个后端(索引器 Indexer),以及它该做什么、不该做什么。
  • 理解后端如何通过 JSON-RPC(HTTP / WebSocket)与节点通信,以及"轮询"和"订阅"两种获取事件方式的成本差异。
  • 看懂事件同步的核心设计:信任边界、确认数、区块重组、两段式推进进度。
  • 用"唯一索引 + 数据库事务"保证数据不重复、不丢失。
  • 跟着一笔真实转账,看它在后端如何被发现、确认、入库、展示。

前置知识

建议先按顺序读完前三篇。本篇默认你已经知道:

  • 合约与 ABI:合约是部署在链上的程序,ABI 是它的"接口说明书"(第 2 篇)。
  • Transfer 事件:MTK 合约每次转账都会发出一条 Transfer 事件(日志),其中 topics 存 from / to,data 存金额(第 2 篇)。
  • JSON-RPC:程序和区块链节点对话的报文格式(第 3 篇)。
  • MetaMask 与 ethers:前端通过钱包签名、发送交易(第 3 篇)。
  • 交易上链过程:签名 → 广播 → 打包进区块 → 生成回执(第 1~3 篇)。

后端是用 Go 写的,你只需要看得懂基本语法;每段代码都会配文字说明。

阅读说明:文中代码摘自本项目源码;为突出主线,部分片段省略了错误处理(标注为"节选"),完整代码以源码为准。 日志与接口返回除特别标注"示意"的以外,均为真实运行结果。


目录

  1. 定位:Web3 应用为什么还需要后端
  2. 技术栈与目录结构
  3. 启动流程:main.go
  4. 配置:config
  5. 数据层:三张表与 models
  6. 与节点通信:JSON-RPC over HTTP / WebSocket
  7. ★ 事件同步:listener 的核心设计
  8. ★ 链上转账:后端的 8 个步骤
  9. 接口层:routers、controllers、统一返回结构
  10. 可靠性:重试、超时、幂等、限流
  11. 动手 Demo:观察后端如何工作
  12. 常见问题
  13. 系列总结

一、定位:Web3 应用为什么还需要后端

这一章要解决的问题:数据明明都在链上,为什么还要写一个后端?它的边界在哪里?

1.1 一句话定位

后端是链上数据的"抄写员"和"翻译员"。 业内把这类服务叫索引器(Indexer):它持续读取链上数据,按业务需要整理成便于查询的形式,就像图书馆给海量藏书编的"索引卡片"。

它不碰钱、不碰私钥,只做两件事:

  1. 抄写:监听 MTK 合约的 Transfer 事件,把每笔转账抄进 MySQL。
  2. 翻译:给地址配上昵称,按"某地址的交易记录"这样的业务需求组织数据,通过 HTTP 接口提供给前端。

1.2 链上已经有数据了,为什么还要抄一份?

链上数据的组织方式是"按区块、按交易",适合验证,不适合按业务查询。对比一下:

需求直接查链查后端数据库
"某地址的全部转账记录,按时间倒序,分页"很难:链上没有这种索引,只能从部署区块开始逐段扫描事件,又慢又耗 RPC 额度一条 SQL,毫秒级
显示昵称链上根本没有昵称users 表
跨链合并展示(anvil + Sepolia)要分别查两条链再合并一张表,network 字段区分
查余额✅ eth_call 免费且权威❌ 不存余额,也不应该存

所以分工是:余额、转账走链;历史记录、昵称走后端。

1.3 后端在整个系统中的位置

┌──────────┐  签名转账(不经过后端)   ┌────────────────────┐
│  前端     │ ───────────────────────▶ │  区块链 MTK 合约     │
│          │                          │  emit Transfer 事件 │
│          │                          └─────────┬──────────┘
│          │                                    │ ① WebSocket 推送(eth_subscribe)
│          │                                    │ ② 确认后 eth_getLogs 查询
│          │                          ┌─────────▼──────────┐
│          │   GET /api/transactions  │  Go 后端             │
│          │ ◀─────────────────────── │  listener + Gin API │
└──────────┘                          └─────────┬──────────┘
                                                │ 写入
                                      ┌─────────▼──────────┐
                                      │  MySQL              │
                                      │  users / transactions│
                                      │  / sync_states       │
                                      └────────────────────┘

注意图中的方向:前端转账不经过后端,直接发到链上;后端只从链上"读",再把整理好的数据提供给前端。

最重要的原则:链上才是真相,数据库只是副本。

  • 数据库可以随时删掉重建:后端会从合约部署区块开始,把所有 Transfer 事件重新抄一遍。
  • 后端停机期间用户照样能转账;重启后会自动补上漏掉的记录。
  • 后端只读链,从不向链发交易,所以不需要私钥,也不花 Gas。

小结:后端 = 索引器。链是"原件",数据库是为了查询方便而做的"副本",副本丢了可以从原件重新抄。


二、技术栈与目录结构

这一章要解决的问题:后端用了哪些工具,代码分别放在哪里。

技术作用
Go 1.27语言
go-ethereum(geth)v1.17与链交互:ethclient 连接节点、订阅事件、查询日志
GinHTTP 路由
GORM + MySQLORM 与数据库
air开发时热重载(make backend)

其中 go-ethereum 是以太坊官方的 Go 实现,它的 ethclient 包相当于 Go 版的 ethers Provider:把 JSON-RPC 调用封装成普通的 Go 函数。ORM(对象关系映射) 让你用结构体操作数据库表,不必手写全部 SQL。

backend/
├── main.go                        # 启动:配置 → 数据库 → 每条链一个 listener → 路由 → HTTP
├── config/config.go               # 读取 .env:多链配置、同步参数
├── models/                        # 数据层
│   ├── core.go                    #   全局 DB + Init():连接 MySQL、建表
│   ├── user.go                    #   User 表 + EnsureUser 幂等创建
│   ├── transaction.go             #   Transaction 表
│   └── sync_state.go              #   SyncState 表(同步进度)
├── listener/                      # ★ 链上事件同步
│   ├── listener.go                #   公共逻辑:连接、推进进度、查询日志、解析入库
│   ├── subscribe.go               #   WebSocket 订阅模式(默认)
│   └── poll.go                    #   HTTP 轮询模式(降级方案)
├── controllers/
│   ├── types.go                   #   统一返回结构 Response{code, msg, data}
│   └── api/                       #   UserController / TransactionController / NetworkController
├── routers/api.go                 # 路由注册
├── middlewares/cors.go            # 跨域
└── .env / .env.example / .air.toml / Dockerfile

整个后端可以分成两半:listener/ 负责"写"(链 → 数据库),controllers/ + routers/ 负责"读"(数据库 → 前端)。本篇的重点是 listener/。


三、启动流程:main.go

这一章要解决的问题:后端启动时按什么顺序做了哪些事。

// main.go(节选)
func main() {
    cfg := config.Load()                                 // 1. 读配置

    if err := models.Init(cfg.MySQLDSN); err != nil {    // 2. 连数据库、建表
        log.Fatal(err)
    }

    ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
    defer stop()

    // 3. 每条链一个 listener goroutine,某条链出错不影响其他链;
    //    节点暂时连不上时 listener 会在后台自动重试,因此启动顺序无所谓
    var listeners []*listener.Listener
    for _, n := range cfg.Networks {
        l := listener.New(n, models.DB, cfg.Batch, cfg.PollInterval, cfg.CheckpointInterval)
        listeners = append(listeners, l)
        go l.Run(ctx)
    }

    router := gin.Default()                              // 4. 注册路由
    router.Use(middlewares.Cors())
    routers.ApiInit(router, listeners)

    srv := &http.Server{Addr: ":" + cfg.Port, Handler: router}
    go func() { srv.ListenAndServe() }()                 // 5. 启动 HTTP 服务

    <-ctx.Done()                                         // 6. 收到 Ctrl+C → 优雅退出
    srv.Shutdown(...)
}

goroutine(Go 协程) 是 Go 里的轻量级线程,go l.Run(ctx) 表示"在后台并行运行这个函数"。

两件事并行运行:listener goroutine 负责"写"(链 → 数据库),HTTP 服务负责"读"(数据库 → 前端)。它们只通过数据库打交道,互不阻塞。

真实启动日志:

[anvil] 开始监听合约 0xDc64a140…F6C9(chainId=31337,确认数=0,方式=WebSocket 订阅)
[anvil] 已通过 WebSocket 订阅 Transfer 事件,区块 2768 之后的事件由节点推送
[sepolia] 开始监听合约 0xCbeE6A18…2447(chainId=11155111,确认数=2,方式=WebSocket 订阅)
HTTP 服务监听 :8080
[sepolia] 已通过 WebSocket 订阅 Transfer 事件,区块 11786453 之后的事件由节点推送

从日志可以看到:两条链各自独立启动,HTTP 服务不需要等 listener 连上节点就能对外服务。"区块 X 之后的事件由节点推送"这句话很关键,第七章会解释它的含义。


四、配置:config

这一章要解决的问题:后端要监听哪条链、哪个合约、从哪个区块开始,这些参数从哪里来。

backend/.env(Key 已隐去):

PORT=8080
MYSQL_DSN=root:@tcp(127.0.0.1:3306)/mini_dapp?charset=utf8mb4&parseTime=True&loc=Local
NETWORKS=anvil,sepolia          # 启用哪些链
BATCH=1000                      # 每次 eth_getLogs 最多查多少个区块
POLL_INTERVAL=12s               # 轮询间隔 / 订阅模式下检查确认数的间隔
CHECKPOINT_INTERVAL=5m          # 订阅模式下,空闲时多久推进一次进度

ANVIL_RPC_URL=http://127.0.0.1:8545
ANVIL_WS_URL=ws://127.0.0.1:8545
ANVIL_CONTRACT=0xdc64a140aa3e981100a9beca4e685f962f0cf6c9
ANVIL_START_BLOCK=517           # 合约部署区块:首次同步从这里开始
ANVIL_CONFIRMATIONS=0           # 本地链不会重组,不用等确认

SEPOLIA_RPC_URL=https://sepolia.infura.io/v3/<API Key>
SEPOLIA_WS_URL=wss://sepolia.infura.io/ws/v3/<API Key>
SEPOLIA_CONTRACT=0xcbee6a188946136b460c833b238eb5665e1f2447
SEPOLIA_START_BLOCK=11784949
SEPOLIA_CONFIRMATIONS=2         # 等 2 个块(约 24 秒)再入库,规避区块重组

几个新词先简单认识一下,第七章会详细讲:

  • START_BLOCK(起始区块):合约部署所在的区块。在它之前合约还不存在,不可能有事件,所以从这里开始扫描。
  • CONFIRMATIONS(确认数):一笔交易所在区块之后又出了几个新块。等的块越多,这笔交易被"撤销"的可能越小。
  • BATCH:一次 eth_getLogs 最多查多少个区块,控制单次请求的范围。

每条链读成一个 NetworkConfig:

type NetworkConfig struct {
    Name          string         // anvil / sepolia,写入 transactions.network
    ChainID       uint64         // 期望的链 ID,启动时与节点返回值比对
    RPCURL        string         // HTTP 地址:普通查询都走它
    WSURL         string         // WebSocket 地址:配置后改用订阅,留空则退回轮询
    Contract      common.Address // MyToken 合约地址
    StartBlock    uint64         // 合约部署区块
    Confirmations uint64         // 确认数
}
配置从哪里来
*_CONTRACT、*_START_BLOCK部署合约后 export-abi.sh 打印
SEPOLIA_RPC_URL、SEPOLIA_WS_URL在 Infura 申请的 API Key 拼出来
ANVIL_*_URLanvil 固定是本机 8545 端口,HTTP 和 WebSocket 共用

五、数据层:三张表与 models

这一章要解决的问题:链上的事件抄到数据库里,要存成什么样子?怎样保证不重复、能续传?

三张表各司其职:users 存"地址 ↔ 昵称",transactions 存转账记录,sync_states 存"抄到哪儿了"。

5.1 users:地址 ↔ 昵称

type User struct {
    ID        uint      `gorm:"primaryKey" json:"id"`
    Address   string    `gorm:"type:char(42);uniqueIndex;not null" json:"address"` // 统一小写
    Nickname  string    `gorm:"type:varchar(32);not null" json:"nickname"`
    CreatedAt time.Time `json:"createdAt"`
}

链上没有"用户"的概念,只有地址。后端把地址映射成用户、分配随机昵称。同一个地址在不同链上是同一个人(地址由私钥推导,与链无关),所以 users 表不区分链。

幂等(Idempotent):同一个操作执行一次和执行多次,结果完全一样。就像电梯的"上楼"按钮,按一次和按十次效果相同。后端大量依赖这个性质,因为重启、重试都可能让同一件事被做多次。

幂等创建用户:

func EnsureUser(db *gorm.DB, address string) (*User, error) {
    addr := strings.ToLower(address)
    if addr == ZeroAddress {
        return nil, nil                  // 零地址(铸造的 from)不是真实用户
    }
    u := User{Address: addr, Nickname: randomNickname()}   // 例如"勇敢的鹦鹉#3717"
    // INSERT ... ON DUPLICATE KEY 什么都不做:地址已存在时不会覆盖原昵称
    if err := db.Clauses(clause.OnConflict{DoNothing: true}).Create(&u).Error; err != nil {
        return nil, err
    }
    var saved User
    db.Where("address = ?", addr).First(&saved)             // 再按地址查出来
    return &saved, nil
}

为什么用"先插入、冲突忽略、再查询",而不是"先查、没有再插"? 因为后者在并发时,两个请求可能同时查到"不存在",然后都去插入。依靠 address 唯一索引,数据库层面保证只会插入一条。

5.2 transactions:Transfer 事件的抄本

type Transaction struct {
    ID       uint   `gorm:"primaryKey"`
    Network  string // anvil / sepolia
    Contract string // 发出事件的合约(同一条链可能部署过多个版本)
    // (chain_id, tx_hash, log_index) 唯一确定一条事件 → 唯一索引保证幂等
    ChainID  uint64 `gorm:"uniqueIndex:uk_chain_tx_log,priority:1"`
    TxHash   string `gorm:"uniqueIndex:uk_chain_tx_log,priority:2"`
    LogIndex uint   `gorm:"uniqueIndex:uk_chain_tx_log,priority:3"`

    FromUserID *uint // 铸造时 from 是零地址,不对应用户,所以可空
    ToUserID   uint

    FromAddr, ToAddr string
    Amount      string    // uint256 最大 78 位,用字符串保存避免精度丢失
    BlockNumber uint64
    BlockTime   time.Time
}

为什么唯一键是 (chain_id, tx_hash, log_index) 三个字段? 一笔交易可能发出多条事件,光靠 tx_hash 区分不开;log_index 是这条事件在区块里的序号;再加上 chain_id,就能在多条链之间也唯一定位一条事件。这相当于给每条事件发了一个"身份证号"。

一条真实记录(Sepolia 上收款人1 转 2500 MTK 给收款人2):

字段值
network / chain_idsepolia / 11155111
tx_hash0x80e152bb556f0d778c136a9461f75cbd080b6a7ebd077a02f2a309f97394f3a5
log_index63
from_user_id → 昵称105 → 勇敢的鹦鹉#3717
to_user_id → 昵称106 → 迷糊的水獭#7322
amount2500000000000000000000(= 2500 MTK)
block_number11786846

5.3 sync_states:同步进度(书签)

type SyncState struct {
    Network   string `gorm:"primaryKey"`
    Contract  string `gorm:"primaryKey"`
    LastBlock uint64 // 已同步到的区块
    UpdatedAt time.Time
}

可以把它想象成读书时夹的书签:链是一本不断变厚的书,每个区块是一页,last_block 记录"已经读完到第几页"。

后端每处理完一段区块,就把书签往后挪。重启后从书签处继续,停机期间的转账不会漏。按 (network, contract) 记录:重新部署合约后,新合约要从它自己的部署区块开始同步。

小结:users 负责"翻译",transactions 负责"抄写",sync_states 负责"记住抄到哪儿"。两个唯一索引(地址、事件三元组)让重复写入自动变成无害操作。


六、与节点通信:JSON-RPC over HTTP / WebSocket

这一章要解决的问题:后端通过什么"语言"和什么"线路"跟区块链节点说话,每种方式的代价是什么。

后端和链之间说的"语言"是 JSON-RPC(报文格式、与 REST 的区别见第 3 篇的「RPC 是什么」和「JSON-RPC、HTTP 与 REST 是什么关系」两节)。和前端不同的是:前端经 MetaMask 转发请求,后端直连 RPC 节点,而且同时用到了 HTTP 和 WebSocket 两种传输方式。

打个比方:JSON-RPC 是"说什么话",HTTP 和 WebSocket 是"用什么线路说"——同样一句话,可以写信(HTTP)寄过去,也可以打一通一直不挂的电话(WebSocket)说。

6.1 后端用到的 JSON-RPC 方法

方法作用go-ethereum 中的调用代码位置
eth_chainId启动时校验节点是不是配置的那条链client.ChainID()listener.go 的 dialAndCheck
eth_blockNumber最新区块号,用来计算安全区块(确认数)client.BlockNumber()advance、subscribeOnce
eth_getLogs按"合约地址 + 事件签名 + 区块范围"查询 Transfer 事件client.FilterLogs()syncRange
eth_getBlockByNumber查区块头,拿区块时间(日志里没有时间)client.HeaderByNumber()blockTime
eth_getTransactionReceipt查交易回执(仅用于回填旧数据的合约地址)client.TransactionReceipt()backfillContract
eth_subscribe订阅 Transfer 事件,由节点主动推送ws.SubscribeFilterLogs()subscribeOnce

后端从不调用 eth_sendRawTransaction:它只读链、不发交易,所以不需要私钥,也不花 Gas。

6.2 HTTP 与 WebSocket:两种传输方式

同一份 JSON-RPC 报文,既可以通过 HTTP 发送,也可以通过 WebSocket 发送。区别在于连接方式:

对比HTTPWebSocket
连接方式每次请求一问一答,答完就结束先用一次 HTTP 请求"升级"成 WebSocket(握手,返回 101 Switching Protocols),之后这条连接一直保持
谁能先说话只能客户端发起双方都可以,节点能主动推送
适合查询订阅新事件(eth_subscribe)
Infura 计费按方法计费建立连接计一次(ws_newConnection),之后按方法 / 推送计费

后端的分工:查询一律走 HTTP(l.client),WebSocket 只用来接收推送(subscribeOnce 里单独建立的 ws)。这样 WebSocket 断线重连时,不会影响正在进行的查询。

6.3 后端用到的节点

节点HTTP 地址(配置项)WebSocket 地址(配置项)
anvil(本地链)http://127.0.0.1:8545(ANVIL_RPC_URL)ws://127.0.0.1:8545(ANVIL_WS_URL),与 HTTP 共用端口
Infura(Sepolia)https://sepolia.infura.io/v3/<Key>(SEPOLIA_RPC_URL)wss://sepolia.infura.io/ws/v3/<Key>(SEPOLIA_WS_URL)

Infura 是一家节点服务商:你不必自己运行以太坊节点,用它提供的地址就能访问链。

Infura 的 Key 写在 backend/.env 里,只存在于服务器上,不会发到用户浏览器。这也是"需要 Key 的节点由后端访问、前端借用 MetaMask 的节点"的原因。

6.4 Infura 计费项 ws_newConnection 是什么

在 Infura 的每日用量统计里,除了各个 JSON-RPC 方法,还会看到一个 ws_newConnection。

它不是 JSON-RPC 方法,而是 Infura 的一个计费项:每新建立一条 WebSocket 连接,按次收一笔费用。 可以这样理解:JSON-RPC 方法是"打电话时说的话",ws_newConnection 是"接通电话"这个动作本身。所以它不放在上面的方法表里,你也无法在请求中写 "method":"ws_newConnection" 去调用它。

本项目中产生这项计费的地方都在后端(前端经 MetaMask 用的是 MetaMask 自己的节点,不消耗你的 Infura 额度):

何时新建 WebSocket 连接代码位置次数
后端启动时校验 WebSocket 地址和链 ID,校验完立即关闭backend/listener/listener.go 的 connect()每次启动 1 次
建立订阅用的长连接backend/listener/subscribe.go 的 subscribeOnce()每次启动 1 次
断线后自动重连同上每次重连 1 次

开发时用 make backend(air 热重载)频繁修改代码,每保存一次后端就会重启一次,这一项就会跟着增加;正常长期运行时,它只在启动和偶尔断线重连时产生。

小结:查询走 HTTP(一问一答),接收推送走 WebSocket(长连接、节点可主动说话)。节点服务按调用计费,所以"少发请求"是后端设计的重要目标——这正是下一章的出发点。


七、★ 事件同步:listener 的核心设计

这一章要解决的问题:怎样又省又稳地把链上的 Transfer 事件抄进数据库——不漏、不重、不写入作废数据,还尽量少花 RPC 额度。

这是全系列概念最密集的一章。先用一个类比建立直觉,再逐个看代码。

7.1 先建立直觉:一个"抄写员"的工作方式

把 listener 想象成一位在档案馆工作的抄写员,任务是把一本不断加页的公告簿(区块链)里所有"MTK 转账公告"(Transfer 事件)抄进自己的笔记本(数据库)。

他会遇到几个现实问题,每个问题对应本章的一个设计:

抄写员遇到的问题对应的技术概念本章小节
要么隔一会儿去翻一遍公告簿,要么请管理员有新公告时打电话通知他轮询(poll)vs 订阅(subscribe)7.2
公告簿里什么公告都有,他只关心 MTK 转账过滤条件(合约地址 + 事件签名)7.3
电话接通之前的旧公告,管理员不会再打电话说,他得自己翻信任边界 trustAfter、补查历史7.4
最新几页偶尔会被撕掉重写,刚贴出来的公告不一定算数区块重组(reorg)、确认数7.5
没接到电话的那些页,能不能不翻、直接把书签往后挪?挪到哪里才安全?advance 两段式推进、subscriptionLag7.6
抄一条公告要同时改好几处,中途被打断怎么办数据库事务、唯一索引幂等7.7

带着这张表往下读,每一节都是在回答其中一个问题。

7.2 两种获取事件的方式:轮询 vs 订阅

  • 轮询(Polling):每隔固定时间主动去问"有没有新事件"。像每隔几分钟去信箱看一眼。
  • 订阅(Subscription):告诉节点"有我关心的事件就通知我",之后由节点主动推送。像让快递员到了直接打电话。
对比HTTP 轮询(poll.go)WebSocket 订阅(subscribe.go,默认)
原理每隔 12 秒问节点"有没有新事件"节点在有匹配事件时主动推送
没有转账时每个块都要 eth_blockNumber + eth_getLogs只有每 5 分钟一次 eth_blockNumber
Sepolia 空闲一天的 Infura 用量(估算)约 240 万 Credits约 2.3 万 Credits
何时使用未配置 *_WS_URL 时自动降级配置了 *_WS_URL

两种方式共用同一套入库逻辑,区别只在"怎么知道有新事件"。

轮询的问题在于:绝大多数时候合约并没有新转账,但每 12 秒仍要问一次,这些请求几乎都是"白问"。订阅把"有没有"这件事交给节点判断,后端平时几乎不用发请求。

真实对比(Infura 后台的每日用量统计):

日期同步方式Infura 用量主要构成
2026-09-26(改造前,统计截至截图时,后端只运行了几个小时)HTTP 轮询约 37.6 万 Creditseth_getLogs 占 74%、eth_blockNumber 占 24%
2026-09-27(改造后,统计截至截图时)WebSocket 订阅5,560 Creditseth_blockNumber 2,800、eth_getLogs 2,040、ws_newConnection 480、eth_call 160、eth_getTransactionReceipt 80

两张截图的统计时长不同,不能直接相除,但数量级的差距一目了然。其中 ws_newConnection 不是 JSON-RPC 方法,而是建立 WebSocket 连接的计费,见 6.4。

7.3 过滤条件:只要 MTK 合约的 Transfer 事件

链上每个区块里有各种合约发出的各种事件。后端通过两个条件把范围缩小到"MTK 合约发出的 Transfer 事件":

  • 合约地址:只要 MTK 合约发出的日志。
  • topics[0]:事件签名的哈希。第 2 篇讲过,日志的第一个 topic 就是 Transfer(address,address,uint256) 的 Keccak256 哈希,用它就能认出"这是一条 Transfer 事件"。
// Transfer 事件签名的哈希 = 日志的 topics[0]
var transferTopic = crypto.Keccak256Hash([]byte("Transfer(address,address,uint256)"))
// = 0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef

func (l *Listener) filterQuery(from, to uint64) ethereum.FilterQuery {
    q := ethereum.FilterQuery{
        Addresses: []common.Address{l.cfg.Contract},   // 只要 MTK 合约发出的
        Topics:    [][]common.Hash{{transferTopic}},    // 只要 Transfer 事件
    }
    if from > 0 || to > 0 {
        q.FromBlock = new(big.Int).SetUint64(from)
        q.ToBlock = new(big.Int).SetUint64(to)
    }
    return q
}

订阅和查询用的是同一个过滤条件。区别只是:订阅时不带区块范围(filterQuery(0, 0),表示"以后的都要");查询时带上 [from, to] 区块范围。

7.4 订阅模式的完整流程

先认识一个关键概念:

信任边界(trustAfter):订阅建立那一刻链上的最新区块号。它把区块分成两部分:

  • 边界之后的新区块:只要有 Transfer 事件,节点一定会推送过来。所以"没收到推送"就等于"没有事件"——这部分区块可以信任订阅。
  • 边界之前的旧区块:订阅不会补推历史事件,节点对它们"只字不提"。所以这部分区块不能信任订阅,必须自己用 eth_getLogs 补查。

回到抄写员的比喻:他和管理员约好"从现在起有新公告就打电话"。挂电话前,他在公告簿上夹一张便签,写着"第 X 页之后靠电话"。第 X 页及之前的内容,他得自己翻。

// listener/subscribe.go(节选:省略了错误处理和超时封装)
func (l *Listener) subscribeOnce(ctx context.Context) (established bool, err error) {
    ws := ethclient.DialContext(ctx, l.cfg.WSURL)              // ① 建立 WebSocket 连接

    logsCh := make(chan types.Log, 128)
    sub := ws.SubscribeFilterLogs(ctx, l.filterQuery(0, 0), logsCh)  // ② eth_subscribe("logs")

    head := l.client.BlockNumber(ctx)
    l.trustAfter = head                                        // ③ 记下"信任边界"
    l.tryAdvance(ctx)                                          // ④ 补查订阅之前的历史区块

    checkpoint := time.NewTicker(l.checkpointInterval)         // 每 5 分钟
    recheck := time.NewTicker(l.pollInterval)                  // 有待确认事件时每 12 秒

    for {
        select {
        case err := <-sub.Err():                               // 连接断了 → 返回,外层重连
            return true, err
        case lg := <-logsCh:                                   // ⑤ 收到推送
            l.onLog(lg)                                        //    记下"区块 X 有事件,等确认"
            if l.cfg.Confirmations == 0 { l.tryAdvance(ctx) }  //    本地链不用等,立即入库
            recheck.Reset(l.pollInterval)
        case <-recheck.C:                                      // ⑥ 检查事件是否已有足够确认
            if l.tryAdvance(ctx) && len(l.pending) == 0 { recheck.Stop() }
        case <-checkpoint.C:                                   // ⑦ 空闲时推进进度
            l.tryAdvance(ctx)
        }
    }
}

逐步说明:

  • ① ②:建立 WebSocket 长连接,并发出 eth_subscribe。从这一刻起,新事件会被推送到 logsCh 通道(channel,Go 里在 goroutine 之间传数据的"管道")。
  • ③:记下当前最新区块作为信任边界。启动日志里的"区块 11786453 之后的事件由节点推送"就是这一步打印的。
  • ④:从书签处开始,用 eth_getLogs 补查到信任边界(首次启动或断线期间的区块)。
  • ⑤ ⑥ ⑦:进入循环,同时等三种信号——收到推送、到了"复查确认"的时间、到了"空闲推进"的时间。select 表示"哪个先来就处理哪个"。

为什么顺序必须是"先订阅 → 记边界 → 再补查"? 如果先补查再订阅,补查结束到订阅建立之间的那几秒里发生的转账,既不在补查范围内,也不会被推送,就会永久漏掉。先订阅就保证了:边界之后的事件一定会被推送,边界之前的由补查负责。

用时间轴看两种顺序的区别:

错误顺序:补查到区块 100 ──(几秒空档:区块 101 的转账没人管)──▶ 订阅建立
正确顺序:订阅建立(边界 = 100,之后的都会推送)──▶ 补查到区块 100

(上图为示意,区块号仅用于说明。)

另外,连接一旦断开,代码会把 trustAfter 恢复成"最大值",意思是"现在没有订阅保障,所有区块都必须真正查询",等重连后再重新设定边界。

7.5 为什么收到推送不直接入库?——确认数与区块重组

// onLog 记录一条推送的事件。这里不直接入库:推送的事件可能还没有足够确认,
// 等确认后统一用 eth_getLogs 按区块范围查询入库,和轮询模式走同一条路径,重组时也不会写入作废的数据。
func (l *Listener) onLog(lg types.Log) {
    if lg.Removed {
        // 区块重组:节点撤回了之前推送的事件。确认后再查询时它自然不会出现
        return
    }
    l.pending[lg.BlockNumber] = struct{}{}
}

区块重组(reorg):公链上偶尔会出现两个验证者同时出块,网络最终只保留其中一条链,另一条链上的交易会被"撤销"。如果收到推送就立即入库,就可能把作废的交易写进数据库。

打个比方:公告栏上最新贴出的几张告示,偶尔会被管理员撕下来换成另一版;越往前的告示,被撕掉的可能性越小。所以抄写员的策略是:最新的几张先不抄,等它后面又贴了几张再抄。

确认数(Confirmations) 就是"后面又贴了几张"的数量。区块 N 之后每多出一个块,就多一个确认。Sepolia 配置为 2,意思是交易所在区块之后再出 2 个块才入库;本地链 anvil 只有一个出块者,不会重组,配置为 0。

所以推送只当作"通知":等这笔交易之后又出了 2 个块(确认数 = 2),再用 eth_getLogs 正式查询。那时查到的就是稳定的结果。

代码里的 pending 就是抄写员的"待办清单":只记下"区块 X 有我的事件",而不是事件本身。等区块 X 足够"老"了,再去正式查一遍。

这样设计还有一个好处:订阅模式和轮询模式最终都通过 eth_getLogs 入库,走同一条代码路径,只需要保证这一条路径正确即可。

7.6 advance:把进度推进到"安全区块"

先定义两个词:

  • head(最新区块):链上当前最新的区块号。
  • safe(安全区块):head - 确认数。编号不超过 safe 的区块已经有足够确认,可以放心入库。

advance 的任务是:把书签从当前位置推进到 safe。这是 listener 最核心的函数,轮询和订阅都调用它:

// listener/listener.go(节选:省略了错误处理)
func (l *Listener) advance(ctx context.Context) error {
    head := l.client.BlockNumber(ctx)               // 链上最新区块
    safe := head - l.cfg.Confirmations              // 安全区块:只处理已有足够确认的
    from := l.nextBlock()                           // 书签 + 1

    // 第 1 段:必须真正 eth_getLogs 查询的区块
    //   - 轮询模式:全部区块(trustAfter = 最大值)
    //   - 订阅模式:订阅建立之前的区块 + 收到过推送的区块
    logsUpTo := min(safe, max(l.trustAfter, l.lastPendingAtOrBelow(safe)))
    for from <= logsUpTo {
        to := min(logsUpTo, from+l.batch-1)         // 每批最多 BATCH 个区块
        l.syncRange(ctx, from, to)                  // eth_getLogs + 入库 + 更新书签
        from = to + 1
    }

    // 第 2 段:订阅保证其中没有事件的区块 → 只更新书签,不花 eth_getLogs
    if trustTo := min(safe, head-subscriptionLag); from <= trustTo {
        saveProgress(l.db, l.cfg.Name, l.contract, trustTo)
    }
    ...
}

为什么要分两段? 因为要推进的区块里,有的"必须亲自查",有的"可以不查直接跳过":

段包含哪些区块怎么处理花不花 eth_getLogs
第 1 段信任边界之前的区块(订阅不负责)+ 收到过推送、且已达到确认数的区块用 eth_getLogs 真正查询、解析、入库花
第 2 段其余区块:在信任边界之后、没收到过推送,说明订阅保证了其中没有 MTK 的 Transfer 事件不查,直接把书签挪过去不花

读 logsUpTo 那一行时可以这样理解:第 1 段的终点 = "信任边界"和"最后一个已确认的推送区块"中较大的那个,但不能超过 safe。

  • 轮询模式没有订阅,trustAfter 是最大值,所以第 1 段就是全部区块,一个都不能跳。
  • 订阅模式下,平时没人转账时第 1 段为空,只有第 2 段在推进书签——这就是订阅模式每 5 分钟只需要一次 eth_blockNumber 的原因。

用数轴来理解(订阅模式,某次 checkpoint 时):

区块号 ──────────────────────────────────────────────────────────────▶
  书签          trustAfter             推送过事件的区块        safe    head
   │   第1段:getLogs  │                     ●                  │       │
   ├──────────────────┤                                        │       │
   │                  │    第2段:订阅保证没事件,直接挪书签       │       │
   │                  ├────────────────────────────────────────┤       │
                                                               └─ 等确认─┘

补充说明图中的 ●:一旦某个推送过事件的区块达到确认数(编号 ≤ safe),lastPendingAtOrBelow 就会返回它,第 1 段会一直延伸到这个区块,保证它被 eth_getLogs 真正查询;还没达到确认数的推送区块先不查,免得为它前面没有事件的区块白白花 eth_getLogs。

subscriptionLag 是什么? 第 2 段的终点不是 safe,而是 min(safe, head - subscriptionLag)。

subscriptionLag = 2:推送从节点传到后端有延迟,所以只把"比最新块至少旧 2 块"的区块当作"确认没有事件",避免推送还在路上时书签就越过了它。

类比:管理员的电话可能晚几秒才打过来。抄写员不能因为"刚才没接到电话"就断定最新一页没有公告——他只对"至少两页之前"的内容下这个结论。对 Sepolia(确认数 2)来说两者相同;对 anvil(确认数 0)来说,这个缓冲才真正起作用。

小结(advance):书签的推进分两步——需要查的区块真查(第 1 段),订阅担保过的区块直接跳(第 2 段);而"担保"只对订阅建立之后、并且足够旧的区块成立。

7.7 syncRange + saveLog:解析事件并入库

advance 决定了"查哪些区块",真正干活的是 syncRange(查询并入库一段区块)和 saveLog(解析并保存一条日志)。

// listener/listener.go(节选)
func (l *Listener) syncRange(ctx context.Context, from, to uint64) error {
    logs := l.client.FilterLogs(ctx, l.filterQuery(from, to))       // eth_getLogs

    // 建用户 + 写交易 + 更新书签放在同一个数据库事务里:要么全部成功,要么全部回滚
    return l.db.Transaction(func(tx *gorm.DB) error {
        for _, lg := range logs {
            l.saveLog(ctx, tx, lg, blockTimes)
        }
        return saveProgress(tx, l.cfg.Name, l.contract, to)
    })
}

数据库事务(Transaction):把多次数据库操作打包成一个整体,要么全部生效,要么全部撤销。就像银行转账"扣款"和"入账"必须同时成功。

为什么要放在一个事务里? 如果"交易已入库但书签没更新",重启后会重复处理(靠唯一索引兜底,问题不大);如果"书签已更新但交易没入库",就会永久漏数据。事务保证两者同进同退。

解析一条日志(对照第 2 篇「★ 链上转账时,合约里发生了什么」一章中的真实日志):

// listener/listener.go(节选:省略了错误处理)
func (l *Listener) saveLog(ctx context.Context, tx *gorm.DB, lg types.Log, ...) (int, error) {
    if lg.Removed || len(lg.Topics) != 3 {
        return 0, nil
    }
    // topics[1] = from,topics[2] = to(indexed 参数,左侧补零到 32 字节,取后 20 字节)
    from := strings.ToLower(common.BytesToAddress(lg.Topics[1].Bytes()).Hex())
    to := strings.ToLower(common.BytesToAddress(lg.Topics[2].Bytes()).Hex())
    // data = value(非 indexed 参数)
    value := new(big.Int).SetBytes(lg.Data)

    blockTime := l.blockTime(ctx, lg.BlockNumber, ...)       // 日志里没有时间,要查区块头

    fromUser := models.EnsureUser(tx, from)                  // 链上地址 → 链下用户(没有就创建)
    toUser := models.EnsureUser(tx, to)

    record := models.Transaction{
        Network: l.cfg.Name, Contract: strings.ToLower(lg.Address.Hex()),
        ChainID: l.cfg.ChainID, TxHash: lg.TxHash.Hex(), LogIndex: lg.Index,
        ToUserID: toUser.ID, FromAddr: from, ToAddr: to,
        Amount: value.String(), BlockNumber: lg.BlockNumber, BlockTime: blockTime,
    }
    // 依靠 (chain_id, tx_hash, log_index) 唯一索引做幂等:重复同步同一区块不会产生重复数据
    res := tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&record)
    return int(res.RowsAffected), res.Error
}

解析过程就是把第 2 篇讲过的日志结构"反过来读":

日志字段真实值(2500 MTK 那笔)解析后写入
topics[1]0x000…00f78c9972…df9f74bf2from_addr = 0xf78c…4bf2 → 勇敢的鹦鹉#3717
topics[2]0x000…00455bddc5…9993d6a0to_addr = 0x455b…d6a0 → 迷糊的水獭#7322
data0x…878678326eac900000amount = "2500000000000000000000"
logIndex0x3flog_index = 63

本章小结:

  1. 订阅负责"及时知道哪里有事件",eth_getLogs 负责"正式、可靠地拿到事件"——推送只是通知,入库一律走查询。
  2. 信任边界之后的区块靠订阅担保,之前的靠补查;顺序必须是"先订阅、再补查"。
  3. 确认数让后端只处理足够"老"的区块,躲开区块重组。
  4. advance 两段式:需要查的真查,订阅担保过的直接挪书签,subscriptionLag 给推送延迟留出余量。
  5. 事务 + 唯一索引:书签和数据同进同退,重复处理也不会产生重复数据。

八、★ 链上转账:后端的 8 个步骤

这一章要解决的问题:把第七章的零件串起来,跟着一笔真实转账走一遍后端。

先回顾一次转账的全流程。下面这张时序图与第 3 篇中的完整时序图是同一张,编号全系列统一:①~⑧ 是前端和钱包签名、广播,⑪~⑭ 是合约执行,⑮ 是前端拿到回执,⑯~⑲ 是本章的主角:后端接力,⑳~㉒ 是前端查询并展示记录。

sequenceDiagram
  autonumber
  actor A as 用户A(发送方)
  participant FE as Vue前端 + ethers
  participant MM as MetaMask
  participant RPC as RPC节点
  participant Chain as 区块链网络(mempool / 出块)
  participant EVM as EVM:MyToken.sol
  participant BE as Go listener
  participant DB as MySQL
  A->>FE: 填写对方地址 B、金额,点击“转账”
  FE->>FE: ABI 编码 calldata = 0xa9059cbb + B + 金额
  FE->>RPC: eth_estimateGas / 获取 nonce 和 Gas 价格
  FE->>MM: 请求签名,交易的 to = 合约地址(不是 B)
  MM->>A: 弹窗:显示 Gas 费并请求确认
  A->>MM: 确认(私钥在本地签名,不外传)
  MM->>RPC: eth_sendRawTransaction(已签名的交易)
  RPC-->>FE: 立即返回 txHash(状态:已广播)
  RPC->>Chain: 校验签名、nonce 和 ETH 余额,放入 mempool 并广播
  Chain->>EVM: 验证者把交易打包进区块,每个节点都执行一遍
  EVM->>EVM: 按函数选择器路由到 transfer,msg.sender = A
  EVM->>EVM: 检查余额,不足则 revert(状态回滚,Gas 照扣)
  EVM->>EVM: SSTORE:_balances[A] 减少,_balances[B] 增加
  EVM->>EVM: emit Transfer(A, B, 金额),写入交易回执的日志
  Chain-->>FE: tx.wait() 拿到回执 status = 1(状态:已上链)
  RPC-->>BE: WebSocket 推送 Transfer 事件(eth_subscribe logs;未配置 WS 时改为定时轮询)
  BE->>RPC: 等够确认数后 eth_getLogs(合约地址, Transfer topic)
  RPC-->>BE: 返回日志:topics = [事件签名, A, B],data = 金额
  BE->>DB: 解码后 INSERT transactions(带 network;按 chain_id + tx_hash + log_index 幂等)
  FE->>BE: GET /api/transactions?address=A,以及 ?address=B
  BE-->>FE: 返回记录,两个列表都显示这笔转账
  FE->>RPC: eth_call balanceOf(A)、balanceOf(B) 刷新余额(免费)

接着第 3 篇前端"交易已上链"(⑮)之后,看后端如何接力(⑯~⑲,以及前端查询到记录的 ⑳~㉑)。以 Sepolia、确认数 2 为例:

步骤 1:收到节点推送

交易被打包进区块 N 的那一刻,Infura 通过 WebSocket 推送这条 Transfer 日志,后端打印(示意):

[sepolia] 收到推送:区块 11786846,交易 0x80e152bb…,等待 2 个确认后入库

pending 记下 11786846,并启动 recheck 定时器。

步骤 2:等待确认

recheck 每 12 秒调用一次 advance:eth_blockNumber 得到 head,计算 safe = head - 2。在 safe < 11786846 之前什么也不做。

步骤 3:确认达成,查询日志

大约 24 秒后,safe >= 11786846。advance 发现有推送过的区块已经安全,于是对 [书签+1, 11786846] 执行 eth_getLogs,拿到正式的日志。

步骤 4:解析日志

saveLog 从 topics 和 data 中解出 from、to、value;再用 eth_getBlockByNumber 查出区块时间。

步骤 5:地址映射成用户

EnsureUser 把两个地址映射成 users 表的记录。收款人第一次出现时,会自动创建并分配一个随机昵称。

步骤 6:写入交易 + 更新书签(同一事务)

后端打印(示意,起始区块取决于当时的书签位置):

[sepolia] 同步区块 11786800→11786846,事件 1 条,新写入 1 条

步骤 7:前端轮询到这条记录

前端的 waitIndexed 每 2 秒调用一次 GET /api/transactions,发现返回里出现了这个 txHash,进度条亮起"后端已入库"。(前端这一侧的代码见第 3 篇「步骤 11:等待后端入库」。)

步骤 8:接口返回(带昵称的记录)

{
  "code": 200,
  "msg": "ok",
  "data": [
    {
      "network": "sepolia",
      "txHash": "0x80e152bb…",
      "logIndex": 63,
      "blockNumber": 11786846,
      "from": "0xf78c997295ae31e35a4e9debaa773d1df9f74bf2",
      "fromNickname": "勇敢的鹦鹉#3717",
      "to": "0x455bddc544682651f75af54a0c2124289993d6a0",
      "toNickname": "迷糊的水獭#7322",
      "amount": "2500000000000000000000",
      "amountDisplay": "2500"
    }
  ]
}

时间线汇总

T+0s    用户在 MetaMask 确认,交易广播
T+12s   交易进入区块 N                  → 前端:打包上链;后端:收到推送(步骤 1)
T+24s   区块 N+1
T+36s   区块 N+2,safe = N              → 后端:查询、解析、入库(步骤 2~6)
T+38s   前端轮询到记录                   → 前端:后端已入库(步骤 7~8)

本地链 anvil 确认数为 0、2 秒一个块,整个过程只要几秒。

小结:从用户看来,"转账成功"和"列表里出现记录"之间有几十秒的间隔,这段时间后端在等确认——这是用一点延迟换取数据可靠。


九、接口层:routers、controllers、统一返回结构

这一章要解决的问题:数据进了数据库之后,前端怎样把它查出来。这部分和普通 Web2 后端基本一样。

9.1 路由

// routers/api.go
func ApiInit(router *gin.Engine, listeners []*listener.Listener) {
    router.GET("/healthz", ...)                                   // K8s 健康检查

    apiRoutes := router.Group("/api")
    {
        apiRoutes.POST("/users", api.UserController{}.Create)
        apiRoutes.GET("/networks", api.NetworkController{Listeners: listeners}.Index)
        apiRoutes.GET("/transactions", api.TransactionController{}.Index)
    }
}

三个接口分别是:登记用户并拿到昵称、查看各条链的同步进度、查询某地址的交易记录。

9.2 统一返回结构

// controllers/types.go
type Response[T any] struct {
    Code int    `json:"code"`   // 与 HTTP 状态码一致,成功为 200
    Msg  string `json:"msg"`
    Data T      `json:"data"`
}

func Success[T any](c *gin.Context, data T) {
    c.JSON(http.StatusOK, Response[T]{Code: http.StatusOK, Msg: "ok", Data: data})
}

func Fail(c *gin.Context, status int, msg string) {
    c.JSON(status, Response[any]{Code: status, Msg: msg, Data: nil})
}

所有接口都返回 {code, msg, data} 这一种形状,前端只需要一套处理逻辑(见第 3 篇「与后端交互:统一请求封装」一章)。

9.3 交易记录查询

// controllers/api/transactionController.go(节选)
func (con TransactionController) Index(c *gin.Context) {
    address := strings.ToLower(c.Query("address"))
    ...
    // 一对多关联查询:该用户转出的 + 转入的;Preload 一次性带出双方用户(昵称)
    q := models.DB.Preload("FromUser").Preload("ToUser").
        Where("(from_user_id = ? OR to_user_id = ?)", user.ID, user.ID)
    if network := c.Query("network"); network != "" {
        q = q.Where("network = ?", network)
    }
    // 只看当前合约:重新部署后,旧合约的记录不会混进来
    if raw := c.Query("contract"); raw != "" {
        q = q.Where("contract IN ?", contracts)
    }
    q.Order("block_number DESC, log_index DESC").Limit(limit).Find(&txs)
    ...
    controllers.Success(c, out)
}

排序用 block_number DESC, log_index DESC:区块号越大越新;同一区块内,按事件序号排。这比按"入库时间"排序更准确,因为补查时入库顺序不等于链上发生的顺序。

金额换算用 decimal,而不是 float64:

// float64 只有约 15~17 位有效数字,1e24 级别的 uint256 会丢精度
func formatAmount(raw string) string {
    v, _ := new(big.Int).SetString(raw, 10)
    return decimal.NewFromBigInt(v, -tokenDecimals).String()   // "2500000000000000000000" → "2500"
}

这与第 3 篇前端用 bigint 处理金额是同一个道理:代币金额以 18 位小数的最小单位存储,数字很大,普通浮点数装不下。

9.4 接口一览(真实返回)

$ curl -s -X POST localhost:8080/api/users -d '{"address":"0x83da6626ef0721da9d452a7eedc4f291a3f9d88c"}'
{"code":200,"msg":"ok","data":{"id":98,"address":"0x83da…d88c","nickname":"合约持有人","createdAt":"…"}}

$ curl -s -X POST localhost:8080/api/users -d '{"address":"abc"}'
{"code":400,"msg":"address 不是合法的以太坊地址","data":null}

$ curl -s localhost:8080/api/networks
{"code":200,"msg":"ok","data":[
  {"network":"anvil","chainId":31337,"contract":"0xdc64…","lastBlock":3860,"latestBlock":3860},
  {"network":"sepolia","chainId":11155111,"contract":"0xcbee…","lastBlock":11786740,"latestBlock":11786742}]}

/api/networks 中的 lastBlock 是书签位置,latestBlock 是链上最新区块。两者的差距反映了确认数和同步延迟:Sepolia 相差 2,正好是配置的确认数。


十、可靠性:重试、超时、幂等、限流

这一章要解决的问题:后端要 7×24 小时和外部节点打交道,网络出问题时怎样不丢数据、不卡死。

一个与外部节点打交道的后台服务,必须假设"网络随时会出问题"。下面这些问题都是本项目开发中真实遇到并修复的:

问题真实现象解决办法
启动时节点连不上anvil 后启动,或 Infura 限流时后端刚好重启 → 这条链被永久跳过,Sepolia 漏掉一笔 8100 MTK 的转账connectWithRetry:连接失败按 1s→2s→…→1min 指数退避重试,直到连上
RPC 请求卡死节点接受连接却不响应 → listener 永久卡住,日志里什么都没有callRPC:每次 RPC 调用最多等 30 秒
限流 429停机 3 小时后重启,补查时连续发了几十次 eth_getLogsBATCH 从 10 调到 1000,一次请求补完;失败后每 12 秒重试
WebSocket 断线网络抖动、Infura 断开空闲连接runSubscribing 自动重连,重连后从书签补查
重复入库重启、重连、重试都可能重复处理同一区块(chain_id, tx_hash, log_index) 唯一索引 + ON CONFLICT DO NOTHING
区块重组公链偶发,已推送的事件被撤回等确认数后再用 eth_getLogs 正式查询

几个术语:

  • 限流(Rate Limit):节点服务商限制单位时间内的请求数,超过后返回 HTTP 状态码 429(Too Many Requests)。
  • 超时(Timeout):给一次请求设定最长等待时间,到点还没回应就放弃,避免无限等待。
  • 指数退避(Exponential Backoff):失败后等待时间逐次翻倍(1 秒、2 秒、4 秒……),既能尽快恢复,又不会在节点故障时疯狂重试把它压垮。

超时封装:

const rpcTimeout = 30 * time.Second

// 超时按"每次调用"单独计算而不是按"整轮同步":补查大量区块总耗时再长也是正常的
func callRPC[T any](ctx context.Context, fn func(context.Context) (T, error)) (T, error) {
    ctx, cancel := context.WithTimeout(ctx, rpcTimeout)
    defer cancel()
    return fn(ctx)
}

// 使用:
head, err := callRPC(ctx, l.client.BlockNumber)

连接重试:

// listener/listener.go(节选:真实代码用 select 监听 ctx,以便 Ctrl+C 时立即退出)
func (l *Listener) connectWithRetry(ctx context.Context) bool {
    backoff := time.Second
    for {
        err := l.connect(ctx)                        // 连 HTTP + WS,并校验链 ID
        if err == nil {
            return true
        }
        if errors.Is(err, errChainIDMismatch) {      // 配置错误(例如把 Sepolia 地址配给了 anvil),重试也没用
            return false
        }
        log.Printf("⚠️  [%s] 连接节点失败,%s 后重试: %v", l.cfg.Name, backoff, err)
        time.Sleep(backoff)
        backoff = min(backoff*2, time.Minute)
    }
}

注意这里区分了两类错误:网络类错误(暂时连不上)值得重试;配置类错误(链 ID 对不上)重试一万次也不会好,应该直接放弃并报错。

小结:可靠性的核心思路是——网络问题靠重试和超时扛过去,重复处理靠幂等兜底,数据不一致靠事务避免,链上的不确定性靠确认数规避。


十一、动手 Demo:观察后端如何工作

这一章要解决的问题:亲手验证前面讲的每个设计。建议按顺序做。

Demo 1:看一笔转账从推送到入库

make anvil          # 终端 1
make backend        # 终端 2,观察日志

在终端 3 用命令行直接转账(不经过页面):

cast send 0xdc64a140aa3e981100a9beca4e685f962f0cf6c9 "transfer(address,uint256)" \
  0x70997970c51812dc3a010c7d01b50e0d17dc79c8 1000000000000000000 \
  --private-key 0xac0974bec39a17e36ba4a6b4d238ff944bacb478cbed5efcae784d7bf4f2ff80 --rpc-url anvil

(这里的私钥是 anvil 内置的公开测试账户私钥,只能用于本地链,切勿在任何真实网络使用。)

终端 2 立刻出现(真实日志):

[anvil] 收到推送:区块 2797,交易 0xc4ad58f8fd3c…,等待 0 个确认后入库
[anvil] 同步区块 2795→2797,事件 1 条,新写入 1 条

这两行分别对应第八章的步骤 1 和步骤 6;因为 anvil 确认数为 0,中间没有等待。

Demo 2:验证幂等——重新同步不会产生重复数据

-- 记下当前记录数
SELECT COUNT(*) FROM transactions WHERE network = 'anvil';

-- 把书签拨回部署区块之前,强制后端重新扫描全部区块
UPDATE sync_states SET last_block = 0 WHERE network = 'anvil';

重启后端,预期日志会显示它重新扫描了全部区块,并且"事件 N 条,新写入 0 条";再查记录数,应该不变。这就是唯一索引 + ON CONFLICT DO NOTHING 的作用(这个实验请你亲自跑一遍验证)。

Demo 3:验证"数据库只是副本"

  1. 停掉后端,在页面上转一笔账:转账照样成功(前端直接和链交互)。
  2. 启动后端:日志显示它补查了停机期间的区块,这笔记录被补进数据库。

Demo 4:切换成轮询模式,对比两种方式

把 backend/.env 里的 SEPOLIA_WS_URL 注释掉,重启后端,日志里的方式变成 方式=HTTP 轮询。之后可以在 Infura 后台的 Stats 页面,对比两种模式下每天的 Credits 消耗。

Demo 5:查看同步进度

SELECT network, contract, last_block, updated_at FROM sync_states;
network  contract                                    last_block  updated_at
anvil    0xdc64a140aa3e981100a9beca4e685f962f0cf6c9  3860        2026-09-26 22:16:16
sepolia  0xcbee6a188946136b460c833b238eb5665e1f2447  11786740    2026-09-26 22:16:20

十二、常见问题

Q:后端为什么不直接从前端接收"我转账了"的通知然后入库? 前端的通知不可信:用户可以伪造请求,说自己转了一万 MTK。后端只相信链上的事件。而且用户可能用别的钱包、别的页面、命令行转账,后端通过监听合约事件就能全部覆盖。

Q:后端需要私钥吗?要付 Gas 吗? 都不需要。后端只读链(订阅、查询),从不发交易。

Q:能不能用后端数据库里的数据算余额? 技术上可以(把所有转入减去转出),但不应该:数据库有延迟、可能漏数据。余额以链上 balanceOf 为准。

Q:后端可以部署多个副本吗? HTTP 接口可以,但 listener 目前是单实例设计:多副本虽然靠唯一索引不会写出重复数据,但会重复消耗 RPC 额度、争抢书签。需要扩容时,应把 listener 拆成独立的单副本服务,或引入选主机制。所以 K8s 部署里后端设为 1 个副本、更新策略用 Recreate。

Q:anvil 重置(不带 --state 重启)后怎么办? 链上数据清空了,但数据库里还有旧记录,书签也比链上最新块还大。后端会提示"已同步到区块 X,但链上最新只有 Y"。执行 make reset-local 清空 anvil 数据,再重新部署合约、更新 .env。


十三、系列总结

作为系列的最后一篇,这一章把四篇文章串起来,回顾一笔转账的完整旅程。

13.1 四篇文章各讲了什么

篇目主题核心收获
第 1 篇Web3 链上转账完整流程(小白版)建立整体认识:账户、代币、合约、链、App 的三个组成部分
第 2 篇智能合约:从零看懂一个代币合约合约是"链上的账本程序":源码、ABI、测试、部署,以及转账时 EVM 如何改账本、发事件
第 3 篇前端:用 MetaMask + ethers 完成链上转账前端通过 RPC 读链、通过钱包签名写链:Provider / Signer、new Contract、转账 12 步
第 4 篇后端:监听链上事件并同步到数据库后端是索引器:订阅事件、等确认、幂等入库,为前端提供链上不便查询的数据

13.2 一笔转账的完整旅程

同一笔转账,在三个部分分别经历了这些步骤:

阶段发生在哪里主要步骤详见
发起前端检测钱包 → 连接钱包 → 切换网络 → 登记用户 → 读余额 → 查询对方 → 表单校验 → 交给钱包签名并广播第 3 篇,步骤 1~8
执行合约(EVM)根据函数选择器找到 transfer → 检查参数和余额 → 修改账本 → 发出 Transfer 事件 → 生成交易回执第 2 篇,第 ⑪~⑮ 步
确认前端等待打包上链 → 刷新双方余额第 3 篇,步骤 9~10
索引后端收到推送 → 等待确认 → 查询日志 → 解析 → 地址映射成用户 → 写入交易并更新书签本篇,步骤 1~6
展示前端 + 后端前端轮询到记录 → 接口返回带昵称的记录 → 交易列表展示第 3 篇步骤 11~12、本篇步骤 7~8

一句话串起来:前端让用户签名并把交易送上链,合约在链上改账本并发出事件,后端听到事件、等它稳定后抄进数据库,前端再把整理好的记录展示给用户。 钱只在链上流动,前端和后端都不"经手"资金。

13.3 进阶学习方向

读完这个系列,如果想继续深入,可以沿着下面几个方向走:

  1. 可升级合约(Upgradeable Contracts):合约部署后代码不可修改,实际项目常用"代理合约 + 实现合约"的模式实现升级。可以从 OpenZeppelin 的代理模式(如 UUPS、Transparent Proxy)入手,重点理解存储布局冲突等风险。
  2. 专业索引服务(如 The Graph):本篇手写的 listener 就是一个最小的索引器。The Graph 等服务把"订阅事件 → 解析 → 存储 → 提供查询接口"做成了通用框架,适合事件种类多、查询复杂的项目。对比它和本项目的实现,能更好地理解索引器的通用问题(重组处理、断点续传、多链)。
  3. 前端工程化(wagmi / TypeChain):wagmi 为 React 提供了钱包连接、读写合约的现成 Hooks;TypeChain 可以根据 ABI 自动生成 TypeScript 类型,替代手写接口,减少"ABI 与代码不一致"的错误。
  4. 后端高可用(多副本 listener 选主):本项目的 listener 是单实例。若要多副本部署,需要引入选主(Leader Election)机制,例如基于数据库锁或 Kubernetes Lease,保证同一时刻只有一个副本在同步,其余副本待命接管。

感谢读完整个系列。希望你现在能从合约、前端、后端三个角度,完整地说清楚"一笔链上转账到底发生了什么"。