基于web3j的以太坊区块数据解析与MySQL持久化实践
简介一套基于 Java Web3j 的以太坊区块数据解析示例工程面向区块链开发者与后端工程师解决直连以太坊节点获取区块数据并存入 MySQL 的核心问题适用于链上数据监控、区块同步分析、区块链教学实验等场景读者具备基础 Java 编程与区块链概念即可上手。资源共 39 个文件压缩包大小 9.64MB主要包含 18 个 jar 依赖库覆盖 Web3j、OkHttp、Jackson、Druid 连接池、MySQL 8.0 驱动以及 RxJava 响应式编程库等另有 6 个 Java 源码与对应 class 文件方便前后对照学习properties 与 config 配置则定义了节点连接地址、数据库连接池和日志输出参数。工程依赖齐备、导入 IDE 即可运行适合作为 Web3j 入门与链上数据二次开发的起点。通过阅读源码可掌握以太坊节点直连、区块高度遍历、交易信息解析、区块数据结构映射等关键逻辑并结合 Druid 连接池将解析结果批量写入 MySQL形成一套可复用的最小实现。目前已有 2082 人学习下载对希望快速搭建区块链数据服务的开发者很有参考价值。1. 直连节点读链上数据比跑全节点轻得多做链上数据分析和区块监控的人早晚会遇到同一个问题Infura 这类托管节点免费额度有限请求频率稍微上来就被限流自己跑一个 geth 全节点几百 GB 的磁盘同步时间又让人劝退。这个基于 java web3j 的小工程走的是第三条路——直连以太坊节点做区块数据解析既可接本地私链也能接远程 RPC解析出的交易、区块、合约数据统一落到 MySQL方便后续查询和二次分析。对需要做链上数据监控、账户流水分析、私链/测试链开发调试的 Java 工程师来说这套代码把 web3j 的调用封装、区块遍历、数据入库串成了一条完整的链路不用再从头摸索。工程依赖的核心组件包括 web3j 5.0.0对应 geth-5.0.0.jar、Druid 连接池、MySQL 8.0 驱动以及 OkHttp 和 RxJava 这些 web3j 的底层依赖。作者把 lib 目录下的 jar 包按传统 Eclipse 工程的方式组织src 下直接放了eth/url.cofig和druid.properties两个配置文件改一下里面的节点地址和数据库连接信息就能跑起来。2. 工程结构与配置加载先看懂这三个类的位置拿到Copy of GethApiCall.zip解压后不要急着去翻代码。这个工程是 Eclipse 直接导出的项目结构.classpath和.settings目录都在用 IntelliJ IDEA 打开时选择 Eclipse 工程导入方式即可。重点看src目录下的三个部分eth包下的是业务代码url.cofig和druid.properties是配置文件lib目录下是全部依赖 jar 包。url.cofig这个文件名容易让人误以为是url.config的拼写错误实际上它就是一个普通的 properties 文件内容大致是节点 RPC 地址和连接参数。常见做法是用Properties类直接读取加载方式如下import java.io.IOException; import java.io.InputStream; import java.util.Properties; public class NodeConfig { private static final Properties props new Properties(); static { try (InputStream in NodeConfig.class.getClassLoader() .getResourceAsStream(eth/url.cofig)) { props.load(in); } catch (IOException e) { throw new ExceptionInInitializerError(加载节点配置失败: e.getMessage()); } } public static String getRpcUrl() { return props.getProperty(eth.rpc.url); } public static int getBlockTimeMs() { return Integer.parseInt(props.getProperty(eth.block.time.ms, 12000)); } }这段代码的核心作用是集中管理节点连接参数。getResourceAsStream(eth/url.cofig)表示配置文件放在 classpath 下的eth目录里如果你的 IDE 没有把src目录标记为资源根目录这里会读取不到运行时会报空指针。将 RPC 地址、区块轮询间隔这些参数外置到 properties 文件里切换主网和测试节点时就不用改代码重新编译直接替换配置即可。Druid 连接池的初始化方式也值得注意。这个工程并没有用 Spring而是在代码里手动创建 DruidDataSource 实例import com.alibaba.druid.pool.DruidDataSource; import java.sql.Connection; import java.sql.SQLException; public class DbManager { private static final DruidDataSource dataSource new DruidDataSource(); static { Properties props new Properties(); try (InputStream in DbManager.class.getClassLoader() .getResourceAsStream(druid.properties)) { props.load(in); dataSource.setUrl(props.getProperty(jdbc.url)); dataSource.setUsername(props.getProperty(jdbc.username)); dataSource.setPassword(props.getProperty(jdbc.password)); dataSource.setDriverClassName(props.getProperty(jdbc.driver)); dataSource.setInitialSize(2); dataSource.setMinIdle(1); dataSource.setMaxActive(10); dataSource.setMaxWait(10000); } catch (IOException e) { throw new ExceptionInInitializerError(加载数据库配置失败); } } public static Connection getConnection() throws SQLException { return dataSource.getConnection(); } }maxWait参数设为 10000 毫秒是关键。如果节点同步慢导致解析线程阻塞数据库连接池的请求线程可能在事务提交前就超时。建议在压测时根据实际吞吐量调大maxActive比如解析速度快但区块间隔短的私链环境连接数 20 比 10 更稳妥。日志层面工程带了log4j.properties和 slf4j-jdk14 的绑定说明日志输出走的是 SLF4J 门面。如果你在 IDEA 里运行时看不到日志检查log4j.properties里的log4j.rootLogger级别是否设为了 DEBUG以及log4j.appender.Console是否指向了org.apache.log4j.ConsoleAppender。3. web3j 直连节点的两种方式与连接池复用用 java web3j 连以太坊节点官方提供了两个入口Web3j.build(HttpService)走 HTTP JSON-RPCWeb3j.build(WebSocketService)走 WebSocket。这个工程用的是 geth-5.0.0.jar对应 web3j 5.x 版本HttpService 是最常用的方式。区别在于 HTTP 每次请求都是独立的 TCP 连接而 WebSocket 是长连接适合需要实时监听区块头或日志的场景。import org.web3j.protocol.Web3j; import org.web3j.protocol.http.HttpService; import org.web3j.protocol.core.methods.response.Web3ClientVersion; public class EthereumNode { private static Web3j web3j; public static synchronized Web3j getInstance() { if (web3j null) { HttpService httpService new HttpService(NodeConfig.getRpcUrl(), new OkHttpClient.Builder() .connectTimeout(10, TimeUnit.SECONDS) .readTimeout(30, TimeUnit.SECONDS) .build()); web3j Web3j.build(httpService); } return web3j; } public static void checkConnection() throws Exception { Web3ClientVersion version web3j.web3ClientVersion().send(); System.out.println(连接成功节点版本: version.getWeb3ClientVersion()); } }web3j实例是线程安全的整个应用只需要一个全局实例不要每次请求都新建。send()方法是同步阻塞调用如果节点响应慢会一直阻塞当前线程所以解析区块的代码必须放在独立线程池里跑。节点地址按环境区分来填环境RPC 地址示例适用场景本地 geth 私链http://127.0.0.1:8545开发调试 / 批量压测远程 geth 全节点http://192.168.1.100:8545生产环境数据采集Erigon / Nethermindhttp://10.0.0.5:8545高性能归档节点一个容易踩的坑是geth 默认只监听127.0.0.1远程访问需要启动时加--http.addr 0.0.0.0 --http.corsdomain *参数。节点启了但连不上时先用curl -X POST http://节点地址:8545 -H Content-Type: application/json -d {jsonrpc:2.0,method:eth_blockNumber,params:[],id:1}测一下端口通不通再用代码连。4. 区块数据解析从区块头到交易明细的字段提取连接建立后核心工作就是区块数据的抓取与解析。web3j 提供了ethGetBlockByNumber方法关键参数有两个区块编号是否使用DefaultBlockParameterName.LATEST以及是否包含完整交易对象。工程里要遍历历史区块做全量解析所以这里用DefaultBlockParameter.valueOf(BigInteger)指定区块号并把returnFullTransactionObjects设为true让 geth 返回完整的交易对象而不是交易哈希列表import org.web3j.protocol.core.DefaultBlockParameter; import org.web3j.protocol.core.methods.response.EthBlock; import org.web3j.protocol.core.methods.response.EthBlock.Block; import org.web3j.protocol.core.methods.response.EthBlock.TransactionObject; import java.math.BigInteger; import java.util.List; public class BlockParser { private final Web3j web3j; public BlockParser(Web3j web3j) { this.web3j web3j; } public ParsedBlock parseBlock(BigInteger blockNumber) throws Exception { EthBlock ethBlock web3j.ethGetBlockByNumber( DefaultBlockParameter.valueOf(blockNumber), true).send(); Block block ethBlock.getBlock(); if (block null) { System.out.println(区块 blockNumber 不存在或尚未同步); return null; } ParsedBlock parsed new ParsedBlock(); parsed.setHeight(block.getNumber()); parsed.setHash(block.getHash()); parsed.setParentHash(block.getParentHash()); parsed.setTimestamp(block.getTimestamp().longValue()); parsed.setTransactionCount(block.getTransactions().size()); parsed.setGasUsed(block.getGasUsed()); parsed.setGasLimit(block.getGasLimit()); parsed.setMiner(block.getMiner()); parsed.setDifficulty(block.getDifficulty()); ListTransactionResult txs block.getTransactions(); for (TransactionResult txResult : txs) { TransactionObject tx (TransactionObject) txResult; parsed.addTransaction(tx); } return parsed; } }getTransactions()返回的是ListTransactionResult当请求参数中returnFullTransactionObjects为false时每个元素是TransactionHash对象只有哈希字符串为true时是TransactionObject包含完整的交易字段。代码里直接强转(TransactionObject)的前提是请求时传了true否则会报ClassCastException。交易对象的字段提取重点处理input和value两项。value的单位是 Wei转成 Ether 要用Convert.fromWei(value, Convert.Unit.ETHER)否则看到的数字是一长串 0input是交易的 calldata普通转账是0x合约调用则是合约方法的 ABI 编码数据import org.web3j.utils.Convert; import org.web3j.utils.Numeric; public class TransactionExtractor { public ParsedTransaction parseTransaction(TransactionObject tx) { ParsedTransaction parsed new ParsedTransaction(); parsed.setTxHash(tx.getHash()); parsed.setFrom(tx.getFrom()); parsed.setTo(tx.getTo() null ? 合约创建 : tx.getTo()); parsed.setValueEther(Convert.fromWei(tx.getValue().toString(), Convert.Unit.ETHER).toPlainString()); parsed.setGasPriceGwei(Convert.fromWei(new BigInteger(tx.getGasPrice().toString()), Convert.Unit.GWEI).toPlainString()); parsed.setGasLimit(tx.getGas()); parsed.setInputData(tx.getInput()); parsed.setNonce(tx.getNonce().longValue()); String input tx.getInput(); if (input ! null input.length() 10 !0x.equals(input)) { String methodId input.substring(0, 10); parsed.setMethodId(methodId); } return parsed; } }getTo()返回null的情况是合约创建交易这个判断不能漏。methodId对应合约方法选择器如果想进一步解析成可读的函数名可以维护一个 methodId 到函数签名的映射表或者调用eth_call做模拟执行。对于普通转账交易input通常是0x不需要做函数签名解析。区块时间戳是 Unix 秒数存入 MySQL 的DATETIME字段前需要转换直接用FROM_UNIXTIME(?)交给 SQL 处理更高效。区块遍历的边界条件是起始区块号取当前已解析的最大高度结束区块号取ethBlockNumber返回值两者相等说明已追到链头进入等待轮询状态。5. MySQL 持久化设计与批量入库优化解析出的区块和交易数据是典型的写多读少场景表结构设计上把区块信息和交易明细拆成两张表通过区块高度关联避免重复存储区块头信息CREATE TABLE block_info ( id bigint(20) NOT NULL AUTO_INCREMENT, block_height bigint(20) NOT NULL COMMENT 区块高度, block_hash varchar(66) NOT NULL COMMENT 区块哈希, parent_hash varchar(66) NOT NULL COMMENT 父区块哈希, block_timestamp datetime NOT NULL COMMENT 出块时间, miner varchar(42) DEFAULT NULL COMMENT 矿工地址, tx_count int(11) NOT NULL DEFAULT 0 COMMENT 交易数量, gas_used bigint(20) DEFAULT NULL, gas_limit bigint(20) DEFAULT NULL, PRIMARY KEY (id), UNIQUE KEY uk_block_height (block_height) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; CREATE TABLE transaction_detail ( id bigint(20) NOT NULL AUTO_INCREMENT, block_height bigint(20) NOT NULL COMMENT 所属区块高度, tx_hash varchar(66) NOT NULL COMMENT 交易哈希, from_address varchar(42) DEFAULT NULL, to_address varchar(42) DEFAULT NULL, value_ether decimal(30,18) DEFAULT NULL COMMENT 转账金额(Ether), gas_price_gwei decimal(30,9) DEFAULT NULL, gas_limit bigint(20) DEFAULT NULL, input_data text COMMENT calldata, method_id varchar(10) DEFAULT NULL, tx_nonce bigint(20) DEFAULT NULL, PRIMARY KEY (id), UNIQUE KEY uk_tx_hash (tx_hash), KEY idx_block_height (block_height), KEY idx_from_address (from_address), KEY idx_to_address (to_address) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;block_hash和tx_hash设置 66 个字符是硬性要求以太坊哈希是 0x 前缀加上 64 位十六进制字符少一位都存不进去。value_ether用decimal(30,18)保存小数因为浮点数在 Web3j 转字符串后再写入会丢失精度。有了uk_block_height唯一索引重跑解析时可以直接用INSERT IGNORE跳过高度的重复入库而不会报错。批量插入用 JDBC 的addBatch()而非逐条执行这是吞吐量提升最明显的一步。解析完一个区块把区块数据加入 batch把该区块下的所有交易也加入 batch攒到 500 条左右统一提交import java.sql.Connection; import java.sql.PreparedStatement; public class BlockDao { private static final String INSERT_BLOCK INSERT INTO block_info (block_height, block_hash, parent_hash, block_timestamp, miner, tx_count, gas_used, gas_limit) VALUES (?, ?, ?, FROM_UNIXTIME(?), ?, ?, ?, ?) ON DUPLICATE KEY UPDATE block_hash VALUES(block_hash); private static final String INSERT_TX INSERT INTO transaction_detail (block_height, tx_hash, from_address, to_address, value_ether, gas_price_gwei, gas_limit, input_data, method_id, tx_nonce) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON DUPLICATE KEY UPDATE tx_hash VALUES(tx_hash); public void batchSave(Connection conn, ParsedBlock block) throws Exception { try (PreparedStatement blockStmt conn.prepareStatement(INSERT_BLOCK); PreparedStatement txStmt conn.prepareStatement(INSERT_TX)) { blockStmt.setLong(1, block.getHeight()); blockStmt.setString(2, block.getHash()); blockStmt.setString(3, block.getParentHash()); blockStmt.setLong(4, block.getTimestamp()); blockStmt.setString(5, block.getMiner()); blockStmt.setInt(6, block.getTransactionCount()); blockStmt.setBigDecimal(7, block.getGasUsed()); blockStmt.setBigDecimal(8, block.getGasLimit()); blockStmt.addBatch(); for (ParsedTransaction tx : block.getTransactions()) { txStmt.setLong(1, block.getHeight()); txStmt.setString(2, tx.getTxHash()); txStmt.setString(3, tx.getFrom()); txStmt.setString(4, tx.getTo()); txStmt.setBigDecimal(5, new BigDecimal(tx.getValueEther())); txStmt.setBigDecimal(6, new BigDecimal(tx.getGasPriceGwei())); txStmt.setLong(7, tx.getGasLimit()); txStmt.setString(8, tx.getInputData()); txStmt.setString(9, tx.getMethodId()); txStmt.setLong(10, tx.getNonce()); txStmt.addBatch(); } blockStmt.executeBatch(); txStmt.executeBatch(); conn.commit(); } } }ON DUPLICATE KEY UPDATE在这里的作用不是真更新而是让重复数据不报错。真正要注意的是conn.setAutoCommit(false)必须打开否则executeBatch()提交的事务会把批量能力抵消掉。传入的Connection从 Druid 连接池取出后统一拿到事务边界再传入 DAO比 DAO 内部自己获取连接更可控。如果启动批量入库后 MySQL 报了PacketTooBigException在 my.cnf 里调大max_allowed_packet即可涉及大批量的 input 数据时检查这个参数。6. 追踪到链头后的等待策略与私链环境批量回补的排错检查区块解析追到链头后进入等待轮询阶段轮询间隔的设置直接影响 RPC 节点压力。以太坊主网平均出块时间是 12 秒但波动很大有的矿池几秒就连出两个块有时候一分钟都没有新块。固定间隔轮询的做法是每 5 秒请求一次ethBlockNumber如果发现最新高度大于本地已入库的最大高度则从后者加 1 开始连续抓取。更平滑的写法是用ethGetBalance之外的方式监听新块——web3j.blockFlowable(false)可以推流式获取但在直连节点场景下轮询比 WebSocket 推送更容易控制节奏import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; public class BlockSyncScheduler { private final Web3j web3j; private final BlockParser parser; private final BlockDao dao; private long syncedHeight; public void start() { ScheduledExecutorService scheduler Executors.newSingleThreadScheduledExecutor(r - { Thread t new Thread(r, block-sync); t.setDaemon(true); return t; }); scheduler.scheduleWithFixedDelay(this::syncOnce, 0, 5, TimeUnit.SECONDS); } private void syncOnce() { try (Connection conn DbManager.getConnection()) { conn.setAutoCommit(false); long latest web3j.ethBlockNumber().send().getBlockNumber().longValue(); if (latest syncedHeight) { for (long h syncedHeight 1; h latest; h) { ParsedBlock block parser.parseBlock(BigInteger.valueOf(h)); if (block ! null) { dao.batchSave(conn, block); syncedHeight h; System.out.println(已同步区块 h 交易数 block.getTransactionCount()); } } } } catch (Exception e) { System.err.println(同步异常等待下一轮重试: e.getMessage()); } } }scheduleWithFixedDelay的意思是上一次任务跑完后间隔 5 秒再开始下一次和scheduleAtFixedRate的固定频率不同。解析耗时长时用前者避免任务堆积解析耗时可忽略时用后者更合适回调频率更均匀。syncedHeight的初始值从数据库SELECT MAX(block_height) FROM block_info恢复避免每次进程重启从头扫链。私链环境里做批量回补测试时常见问题集中在三处。第一geth 私链如果没开启--http.api eth,net,web3中的eth模块应用会直接返回method not found加上--http.api eth,web3,net后重启节点才行。第二默认 geth 最大同时请求数是 500批量回补时的并发请求超过该限制会出现request limit exceeded改启动参数加--rpc.gascap 0和--ws.origins的同时在代码里把HttpService的调用改成单线程串行即可。第三parseBlock里的EthBlock ethBlock web3j.ethGetBlockByNumber(...).send()如果send()连续超时确认 OkHttp 的readTimeout是否大于节点实际响应时间——私链节点在同时处理大量解析请求时单次区块响应可能超过默认 30 秒。本文还有配套的精品资源点击获取