SeaTunnel Web3j 源连接器:通过 JSON-RPC 接入区块链区块高度数据的完整指南
2026/9/17 8:07:41 网站建设 项目流程

SeaTunnel Web3j 源连接器:通过 JSON-RPC 接入区块链区块高度数据的完整指南

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

Web3j Source 是 SeaTunnel 提供的区块链数据源连接器,通过 Web3 provider 的 JSON-RPC 端点读取以太坊网络的最新区块高度,并以单行 JSON 数据的形式向下游输出。读完本文,你将了解该连接器的功能边界(批/流模式、不支持并行与 exactly-once)、唯一配置项url的校验规则、Batch 与 Streaming 两种模式下的完整作业配置,以及底层基于 web3j 4.8.4 的HttpService调用链路与源码实现细节。

功能定位与支持能力

Web3j 连接器是一个只读源(Source)连接器,它通过 Web3 provider 端点读取区块链数据。当前实现只读取最新区块号(latest block number),并输出一个名为value的字段,该字段是包含blockNumber和读取时间戳的 JSON 字符串。

在两种作业模式下其行为如下:

  • 批模式(BATCH):源只发出一行数据即结束;
  • 流模式(STREAMING):持续轮询 provider,每轮输出当前观测到的最新区块号。

从 功能说明 给出的特性矩阵看,该连接器的能力边界非常明确:

能力是否支持
batch支持
stream支持
exactly-once不支持
column projection不支持
parallelism不支持
user-defined split不支持

连接器使用单个 split,不支持并行。由于每行数据对应一次 HTTPeth_blockNumber调用,实际的轮询节奏由所配置 provider 的响应时间决定。

它声明支持 Spark、Flink、SeaTunnel Zeta 三类引擎,属于 SeaTunnel Connector V2 体系下的标准源连接器。

输出 Schema 与数据载荷

连接器只暴露单行表类型,schema 定义如下:

字段类型说明
valueStringJSON 字符串,包含最新区块号以及由连接器生成的时间戳

value字段中 JSON 的结构为:

{"blockNumber":19525949,"timestamp":"2024-03-27T13:28:45.605Z"}
  • blockNumber是 provider 返回的最新区块头高度(数值型);
  • timestamp是连接器在观测到响应的那一刻用Instant.now()生成的 ISO-8601 UTC 时间戳。

连接器不会解析载荷中的其他内容。若你需要交易数、gas 用量、节点元数据等额外字段,需要调用其他 RPC 或在上游用 SQL Transform 做后处理(见文末 FAQ)。

配置项:唯一的url及其校验规则

参数说明

名称类型是否必填默认值说明
urlString-与 Ethereum 网络通信的 Web3 provider 端点,例如 Infura URL

必填的url不允许为空或仅包含空白字符。

源码中的校验实现

配置项在 Web3jSourceOptions.java 中定义:

public static final Option<String> URL = Options.key("url") .stringType() .noDefaultValue() .withDescription( "your infura project url like : https://mainnet.infura.io/v3/xxxxxxxxxxxx");

可以看到该选项声明为字符串类型、无默认值。校验规则在 Web3jSourceFactory.java 中通过OptionRule声明:

@Override public OptionRule optionRule() { return OptionRule.builder().required(URL, notBlank(URL)).build(); }

url必须存在,且满足notBlank条件(不能是空白串)。工厂类通过@AutoService(Factory.class)注解注册到 SeaTunnel 的插件发现机制中,插件标识符(factoryIdentifier())为Web3j——这就是作业配置中source { Web3j { ... } }的由来。

这些校验规则由单元测试 Web3jFactoryTest.java 完整覆盖:

  • testNonblankUrlAcceptedhttp://localhost:8545https://example.com、甚至带首尾空格的custom-endpoint都能通过校验;
  • testMissingUrlRejected/testEmptyUrlRejected/testWhitespaceOnlyUrlRejected:缺失url、空字符串""、纯空白" ""\t\r\n"都会抛出OptionValidationException

这证实了文档中"url 必须非空且非纯空白"的说法,也与notBlank(URL)规则一一对应。

Batch 模式配置示例

在批模式下,源发出一行数据后作业即完成。完整配置如下:

env { parallelism = 1 job.mode = "BATCH" } source { Web3j { url = "https://mainnet.infura.io/v3/xxxxx" plugin_output = "web3j" } } sink { Console { plugin_input = "web3j" parallelism = 1 } }

运行后可观察到如下数据(value字段为转义后的 JSON 字符串):

{"value":"{\"blockNumber\":19525949,\"timestamp\":\"2024-03-27T13:28:45.605Z\"}"}

Streaming 模式配置示例

在流模式下,连接器持续轮询 provider,每轮输出一行包含最新区块号的数据。官方文档给出的示例使用 Assert 连接器校验value非空:

env { parallelism = 1 job.mode = "STREAMING" checkpoint.interval = 10000 } source { Web3j { url = "https://mainnet.infura.io/v3/xxxxx" plugin_output = "web3j" } } sink { Assert { plugin_input = "web3j" rules { field_rules = [ { field_name = value field_type = string field_value = [ { rule_type = NOT_NULL } ] } ] } } }

注意checkpoint.interval = 10000:由于源本身不支持 exactly-once,checkpoint 仅用于下游 sink 需要时的可重放背压控制(详见 FAQ)。

源码级实现剖析

Web3j 连接器模块位于 seatunnel-connectors-v2/connector-web3j,代码量很小,核心只有 4 个类 + 1 个测试,非常适合通读。其依赖方面,pom.xml 声明了 web3jcore4.8.4版本,以及 SeaTunnel 内部的connector-commonseatunnel-format-json模块。

类结构与调用链

Web3jSourceFactory (插件发现入口, identifier="Web3j") └──> Web3jSource (AbstractSingleSplitSource 子类) ├── Web3jSourceParameter (只持有 url) └── Web3jSourceReader (AbstractSingleSplitReader 子类) └── org.web3j.protocol.Web3j + HttpService

1. 有界性由作业模式决定—— Web3jSource.java 中:

@Override public Boundedness getBoundedness() { return JobMode.BATCH.equals(jobContext.getJobMode()) ? Boundedness.BOUNDED : Boundedness.UNBOUNDED; }

即 Batch 模式为BOUNDED、流模式为UNBOUNDED,这正是"批模式发一行即结束、流模式无限轮询"行为的根源。同时getProducedCatalogTables()声明了输出表:表标识为Web3j,schema 只有一个value列,类型为BasicType.STRING_TYPE——与文档的 Output Schema 完全一致。

2. 单 split 单 reader 的轮询逻辑—— Web3jSourceReader.java 是核心:

  • open()阶段用Web3j.build(new HttpService(url))建立与 provider 的 HTTP JSON-RPC 连接;
  • close()阶段调用web3.shutdown()释放连接;
  • pollNext()阶段发起一次轮询:
web3.ethBlockNumber() .flowable() .subscribe( blockNumber -> { Map<String, Object> data = new HashMap<>(); data.put("timestamp", Instant.now().toString()); data.put("blockNumber", blockNumber.getBlockNumber()); String json = OBJECT_MAPPER.writeValueAsString(data); output.collect(new SeaTunnelRow(new Object[] {json})); if (Boundedness.BOUNDED.equals(context.getBoundedness())) { // signal to the source that we have reached the end of the data. context.signalNoMoreElement(); } });

几个关键点值得注意:

  • 每次pollNext通过 web3j 的ethBlockNumber().flowable().subscribe(...)发起一次eth_blockNumberRPC 调用;
  • 响应到达后,reader 组装timestamp(本地生成)与blockNumber(provider 返回),用静态共享的ObjectMapper序列化为 JSON 字符串,打包成单列SeaTunnelRow向下游收集;
  • 只有在BOUNDED(即批模式)时才会调用context.signalNoMoreElement()通知框架数据已尽——这就是批模式"一行即止"的实现机制;流模式下该信号永远不触发,源持续运行。

3. 参数传递—— Web3jSourceParameter.java 只是一个可序列化的包装类,从ReadonlyConfig中取出url并暴露给 reader 使用,没有任何额外配置项——这与"唯一配置项是url"的文档描述吻合。

从源码结构看,轮询是同步阻塞式的:pollNext发起调用并等待 provider 返回后才收集一行,因此不存在独立的轮询间隔配置,实际产出节奏等于"RPC 往返时间"的倒数。

使用注意事项(Notes)

文档明确列出以下注意事项,结合源码均可验证:

  • url必须指向JSON-RPC 兼容的 Web3 provider,如 Infura、Alchemy 或自托管的 Ethereum 节点。推荐 HTTPS;连接器不做额外鉴权,若 provider 要求 API key,请直接把 key 写进 URL(如https://mainnet.infura.io/v3/xxxxx)。
  • 连接器只暴露含value字段的单行类型。提取blockNumbertimestamp用于后续处理,需借助下游的 SQL Transform 或 JSON path。
  • 流模式下连接器保持 HTTP 连接打开,每轮输出观测到的最新区块号;只有当下游 sink 需要时才配合 checkpoint 使用。

FAQ

Q1:轮询速率如何控制?

连接器每次轮询发起一次 HTTPeth_blockNumberRPC,阻塞等待 provider 返回,然后把结果作为一行输出。因此实际产出节奏完全取决于所配置 provider 的响应时间——它无法通过任何配置项独立设置。如果你需要可重放的背压控制,请搭配使用 checkpoint 的下游 sink;Web3j 源本身没有速率限制或节流设置。

Q2:value字段包含什么载荷?

一个 schema 为{"blockNumber": <number>, "timestamp": "<ISO-8601 UTC>"}的 JSON 对象。blockNumber是 provider 返回的最新区块头,timestamp是连接器在观测到响应那一刻生成的。连接器从不检查载荷内容;若需要额外字段(交易数、gas 用量、节点元数据等),必须调用其他 RPC 或在上游使用 SQL Transform 做后处理。

Q3:Web3j 源是否支持 URL 之外的鉴权方式?

不支持。连接器把配置的url原样传给底层 HTTP 客户端。因此鉴权方式就是 provider 在 URL 中接受的形式——通常是 Infura、Alchemy 等托管服务内嵌的 API key。没有单独的基于 Header 的鉴权路径,所以不要把需要频繁轮换的密钥放进 URL;应签发长期有效的 provider key,并通过常规的密钥管理系统管理。

变更记录

变更说明版本
[improve] 更新 Web3j 连接器配置选项(#9005)调整配置项(引入url校验等)2.3.10
[Feature][Connector-V2] 新增 web3j source 连接器(#6598)连接器首次引入2.3.6

完整变更记录见 connector-web3j.md。

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询