You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Web3j获取智能合约事件失败排查及可运行代码示例求助

Web3j监听智能合约事件问题排查与可运行示例

问题根因说明

  • 你当前代码核心问题是EthFilter同时使用DefaultBlockParameterName.LATEST作为起始块和结束块时,绝大多数以太坊RPC节点仅会返回过滤器创建时刻最新高度区块的日志,不会自动订阅后续新产生的区块日志
  • 忽略错误日志打印,RPC节点返回的异常会被RxJava流默认吞掉,无法定位问题
  • 事件定义如果和合约ABI不匹配(比如indexed参数标记错误、事件名/参数类型顺序错),生成的topic会完全匹配不到对应事件
  • 部分RPC节点对大小写敏感的合约地址过滤会失效,需要统一地址格式

可运行完整代码示例

首先引入稳定版Web3j依赖(Maven为例):

<dependency>
    <groupId>org.web3j</groupId>
    <artifactId>core</artifactId>
    <version>4.9.8</version>
</dependency>

完整实现代码:

import org.web3j.protocol.Web3j;
import org.web3j.protocol.core.DefaultBlockParameter;
import org.web3j.protocol.core.DefaultBlockParameterName;
import org.web3j.protocol.core.methods.request.EthFilter;
import org.web3j.protocol.core.methods.response.Log;
import org.web3j.protocol.http.HttpService;
import org.web3j.abi.EventEncoder;
import org.web3j.abi.TypeReference;
import org.web3j.abi.datatypes.Address;
import org.web3j.abi.datatypes.Event;
import org.web3j.abi.datatypes.generated.Uint256;
import io.reactivex.disposables.Disposable;
import java.util.Arrays;
import java.util.concurrent.CopyOnWriteArrayList;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public class ContractEventSubscriber {
    private static final Logger log = LoggerFactory.getLogger(ContractEventSubscriber.class);
    private final Web3j web3j;
    private final CopyOnWriteArrayList<Disposable> subscriptions = new CopyOnWriteArrayList<>();
    private int logCount = 0;

    // 初始化Web3j客户端,替换为你自己的RPC节点地址
    public ContractEventSubscriber() {
        this.web3j = Web3j.build(new HttpService("https://你的RPC节点地址"));
    }

    // 日志消费类
    class LogCounter implements io.reactivex.functions.Consumer<Log> {
        @Override
        public void accept(Log logObject) throws Exception {
            logCount++;
            log.info("收到事件 #{} 区块高度:{} 交易哈希:{} 日志内容:{}", 
                logCount, logObject.getBlockNumber(), logObject.getTransactionHash(), logObject);
        }
    }

    /**
     * 持续订阅新产生的合约事件
     * @param contractAddress 目标合约地址
     */
    public void subscribeNewEvents(String contractAddress) {
        log.info("开始订阅合约{}的事件", contractAddress);
        EthFilter filter = new EthFilter(
                DefaultBlockParameterName.LATEST,
                DefaultBlockParameterName.LATEST,
                contractAddress.toLowerCase() // 统一转为小写避免大小写匹配问题
        );

        // 配置要监听的事件,确保和合约ABI完全一致
        Event withdrawEvent = new Event("Withdraw", Arrays.asList(
                new TypeReference<Address>(true) {}, // 合约中indexed的参数要标记true,否则topic匹配不上
                new TypeReference<Uint256>(false) {},
                new TypeReference<Uint256>(false) {}
        ));
        filter.addSingleTopic(EventEncoder.encode(withdrawEvent));
        // 不需要过滤特定事件可以删掉上面两行配置,就能拉取该合约所有事件

        LogCounter logCounter = new LogCounter();
        Disposable subscription = web3j
                .ethLogFlowable(filter)
                .retry() // 网络波动自动重试
                .subscribe(
                        logCounter,
                        error -> log.error("事件监听出错", error), // 必须加错误打印,否则异常会被吞
                        () -> log.info("事件监听流结束")
                );
        subscriptions.add(subscription);
        log.info("合约事件订阅完成");
    }

    /**
     * 回溯拉取历史区块的合约事件,解决RPC节点最大块范围限制问题
     * @param contractAddress 目标合约地址
     * @param startBlock 起始区块高度
     * @param endBlock 结束区块高度
     */
    public void pullHistoricalEvents(String contractAddress, long startBlock, long endBlock) throws Exception {
        // 按每次拉取4000个块拆分请求,避开节点5000块的范围限制
        long step = 4000;
        for (long currentStart = startBlock; currentStart <= endBlock; currentStart += step) {
            long currentEnd = Math.min(currentStart + step - 1, endBlock);
            log.info("拉取区块{}~{}的事件", currentStart, currentEnd);

            EthFilter historyFilter = new EthFilter(
                    DefaultBlockParameter.valueOf(currentStart),
                    DefaultBlockParameter.valueOf(currentEnd),
                    contractAddress.toLowerCase()
            );
            // 事件过滤规则和订阅逻辑保持一致
            Event withdrawEvent = new Event("Withdraw", Arrays.asList(
                    new TypeReference<Address>(true) {},
                    new TypeReference<Uint256>(false) {},
                    new TypeReference<Uint256>(false) {}
            ));
            historyFilter.addSingleTopic(EventEncoder.encode(withdrawEvent));

            // 批量拉取历史日志
            web3j.ethGetLogs(historyFilter).send().getLogs().forEach(logObj -> {
                Log logObject = (Log) logObj.get();
                logCount++;
                log.info("历史事件 #{} 区块高度:{} 交易哈希:{}", 
                    logCount, logObject.getBlockNumber(), logObject.getTransactionHash());
            });

            // 加延迟避免触发RPC节点频率限制
            Thread.sleep(1000);
        }
    }

    // 销毁时释放资源
    public void shutdown() {
        subscriptions.forEach(Disposable::dispose);
        web3j.shutdown();
    }

    public static void main(String[] args) throws Exception {
        ContractEventSubscriber subscriber = new ContractEventSubscriber();
        // 替换为你的合约地址
        String contractAddr = "0x你的合约地址";
        // 订阅新事件
        subscriber.subscribeNewEvents(contractAddr);
        // 需要拉取历史事件可调用下方方法,替换为目标区块范围
        // subscriber.pullHistoricalEvents(contractAddr, 18000000L, 18050000L);

        // 保持进程运行
        Thread.currentThread().join();
    }
}

额外优化建议

  • 如果使用Infura、Alchemy这类公共RPC节点,推荐用websocket接口替代http接口,事件推送延迟更低、稳定性更好,只需将HttpService替换为WebSocketService即可
  • 如果事件监听延迟较高,可自行调整Web3j的轮询间隔,默认间隔为15秒

内容的提问来源于stack exchange,提问作者Paul Verest on LinkedIn

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.05 19:42:03