SeaTunnel 首个任务实战:基于 FakeSource 与 Console 的本地全链路验证指南
2026/9/18 8:13:14 网站建设 项目流程

SeaTunnel 首个任务实战:基于 FakeSource 与 Console 的本地全链路验证指南

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

本篇指南带你走完 Apache SeaTunnel 的"最短成功路径":全程在本机完成,不依赖 MySQL、Kafka 或任何对象存储,用内置的FakeSource生成数据、FieldMapper做字段重命名、Console输出到终端,一次性验证安装、配置解析与执行引擎三个关键环节是否正常工作。读完本文,你将掌握插件裁剪安装、HOCON 任务配置编写、-m local本地模式提交以及结果判定的完整方法,可以放心进入真实管道开发。

Step 1:完成本地部署

运行本文示例的前提是已经完成 SeaTunnel 的本地部署,并确认 SeaTunnel 安装目录(下文统一称${SEATUNNEL_HOME})下存在可执行的bin/seatunnel.sh。完整的部署流程请参见 本地部署文档,其要点如下:

  • 环境依赖:安装 Java 8 或 11(高于 Java 8 的版本理论上也可用),并配置好JAVA_HOME
  • 获取发行包:从官方下载页获取seatunnel-<version>-bin.tar.gz二进制包并解压(Windows 使用对应的.zip包);
  • 确认脚本就绪:解压后目录内应包含bin/seatunnel.shbin/seatunnel.cmdconfig/connectors/等目录。

在仓库中,bin/目录实际提供了两个安装脚本:install-plugin.sh 与install-plugin.cmd;发行版中的seatunnel.sh会在构建分发时一并产出(源码侧对应的提交入口可参见 SeaTunnelClient.java)。

Step 2:只安装示例所需的插件

从 2.2.0-beta 版本起,官方二进制包默认不再附带连接器依赖,首次使用前必须执行插件安装命令。而生产环境通常也不需要全部插件,因此推荐的做法是:先在config/plugin_config中声明本次任务真正需要的插件,再执行安装。

精简 plugin_config

按 部署文档 > Download The Connector Plugins 的说明,将 config/plugin_config 精简为仅保留connector-fakeconnector-console两项:

--seatunnel-connectors-- connector-fake connector-console --end--

关于plugin_config需要说明两点(均可在仓库中直接核对):

  • 仓库中真实的 config/plugin_config 使用--connectors-v2--作为段标记,并包含了全部已注册连接器的 artifactId(如connector-jdbcconnector-kafkaconnector-cdc-mysql等),文件头部注释明确写道:不要修改分隔符--,只需挑选你需要的插件
  • 完整的插件清单与 artifactId 映射还可通过发行包内的connectors/plugins-mapping.properties(以及仓库根目录的 plugin-mapping.properties)查看,脚本安装时正是以config/plugin_config中列出的行为准逐个下载。

执行安装并核对结果

cd "${SEATUNNEL_HOME}" sh bin/install-plugin.sh ls connectors | rg 'connector-(fake|console)'

仓库中的 bin/install-plugin.sh 展示了这套脚本的实际行为,了解它有助于排查安装问题:

  • 默认插件版本固定为3.0.0,也支持通过第一个参数指定版本:sh bin/install-plugin.sh 3.0.0
  • 默认走 HTTPS 直连下载(需要curlmktemp以及sha512sum/sha1sum/shasum/openssl之一做校验和验证),下载后还会校验 JAR 魔数(504b)确保文件确实是 ZIP/JAR 格式;
  • 可通过环境变量调整行为:SEATUNNEL_MAVEN_REPOSITORY指定 HTTPS Maven 兼容镜像地址;SEATUNNEL_PLUGIN_DOWNLOAD_METHOD=maven则改用项目自带的 Maven Wrapper(mvnw dependency:get)下载,从而支持镜像、私有仓库、代理等 Mavensettings.xml能力;
  • 对于SNAPSHOTLATESTRELEASE及版本区间等无法直接解析的动态版本,脚本会自动切换为 Maven 方式下载。

若脚本运行后connectors/目录下能看到connector-fake-3.0.0.jarconnector-console-3.0.0.jar,则插件安装成功。Windows 用户请使用bin\install-plugin.cmd(该脚本使用捆绑的 Maven Wrapper,无需单独安装 Maven)。

Step 3:编写一个最小任务配置

将下面的配置保存为config/v2.batch.config.template或任意本地文件。它是标准的 SeaTunnel HOCON 配置,由envsourcetransformsink四大块组成:

env { parallelism = 1 job.mode = "BATCH" } source { FakeSource { plugin_output = "fake" row.num = 16 schema = { fields { name = "string" age = "int" } } } } transform { FieldMapper { plugin_input = "fake" plugin_output = "fake1" field_mapper = { age = age name = new_name } } } sink { Console { plugin_input = "fake1" } }

对照仓库中的 config/v2.batch.config.template 可以看出,本示例在此基础上做了三个关键改动:将parallelism降为1(去掉checkpoint.interval)、引入FieldMapper变换并显式打通plugin_output/plugin_input的数据流命名,便于你理解 SeaTunnel 的插件数据通道(channel)机制。

逐块解读配置

env 块parallelism = 1表示整个作业以单并行度运行,便于观察输出顺序;job.mode = "BATCH"声明这是批处理模式(对应示例模板中默认即为 BATCH)。如果你的环境模板里还带有checkpoint.interval,在本示例中并非必需。

source 块(FakeSource)FakeSource是一个纯内存的数据生成器,专门用于测试与演示,不需要任何外部依赖。FakeSourceOptions.java 中定义了它的全部可选参数,常用的包括:

参数默认值说明
row.num5每个并行度生成的数据总行数(示例显式设为 16)
split.num1每个并行度由 enumerator 切分的 split 数量
split.read-interval1reader 两次 split 读取之间的间隔(毫秒)
string.length5生成的 string 类型字段长度
map.size/array.size/bytes.length5对应复杂类型的生成尺寸
string.template若配置,则 string 字段从模板列表中随机选取
int.min/int.max0/Integer.MAX_VALUE整型生成范围,其他数值类型同理(tinyintsmallintbigintfloatdouble等均有min/max对)
rows显式指定要输出的行列表(每行含kindfields),优先级高于随机生成
string.fake.modeRANGE生成模式,可选RANGE(区间随机)或TEMPLATE(模板选取)
auto.increment.enabledfalse是否启用自增 ID 生成,配合auto.increment.start(默认 1)使用

row.num的"每并行度"语义很重要:若你设置parallelism = 2row.num = 16,实际会生成 32 行数据,因为 FakeSourceOptions.java 中明确描述其为 "The total number of data generated per degree of parallelism"。此外schema.fields定义了输出表的字段结构(此处为name: stringage: int),它最终会被解析为CatalogTable(见 FakeConfig.java)。

transform 块(FieldMapper)FieldMapper负责输入输出字段的映射与重命名。field_mapper是一个"源字段 -> 目标字段"的映射表,其配置定义见 FieldMapperTransformConfig.java。示例中age = age表示原样保留age字段,name = new_name表示将name字段重命名为new_name。通过plugin_input = "fake"plugin_output = "fake1"将上游FakeSource的输出通道fake接入、并把变换后的结果输出到新通道fake1

sink 块(Console)Console是一个"打印到终端"的调试型 Sink,同样无需任何外部系统,通过plugin_input = "fake1"消费变换后的数据流。

Step 4:以本地模式运行

进入解压后的 SeaTunnel 目录,使用-m local指定本地模式提交任务:

cd "apache-seatunnel-${version}" ./bin/seatunnel.sh --config ./config/v2.batch.config.template -m local

命令参数的解析实现在 ClientCommandArgs.java 中,-m--master,别名-e/--deploy-mode)的取值被限定为localcluster两种:-m local表示由客户端本地拉起引擎执行,无需预先部署 SeaTunnel Engine 服务集群,这也是新用户最推荐的入门方式;cluster模式则要求先部署引擎服务。-e/--deploy-mode自 2.3.1 起已被标记为 deprecated,建议统一使用-m/--master

--config-m外,该参数类还提供了若干实用选项,可用于后续调试:

  • -d/--dry-run:在不真正运行 Sink 的前提下做校验或预览,支持static(仅静态校验配置)、connect(校验连接)、sample(采样预览数据)三种模式;
  • --sample-limitsample模式下每个 Source 最多转发行数(默认 10,上限 10000);
  • -l/--list:列出作业状态;-j/--job-id:按 JobId 查询作业状态。

预期验证结果

任务提交后,应按以下四条标准逐项核对(全部满足即说明本机基础链路健康):

  1. 进程正常启动,无连接器加载错误:说明connector-fakeconnector-console两个 JAR 被正确发现并加载;
  2. 控制台打印output rowType:这一行由 ConsoleSinkWriter.java 在初始化时以log.info("output rowType: {}", ...)输出,内容应展示经过FieldMapper映射后的字段,即agenew_name,这是验证 transform 生效的直接证据;
  3. 控制台打印来自ConsoleSinkWriter的 16 行数据:每条记录形如subtaskIndex=0 rowIndex=N: SeaTunnelRow#tableId=... SeaTunnelRow#kind=INSERT : <字段值>(见 ConsoleSinkWriter.java),这里出现 16 行正是因为parallelism = 1row.num = 16
  4. 批任务在写完所有行后正常退出BATCH模式下数据有界,作业完成后客户端应正常结束并返回。

如果你修改了parallelismrow.num,请记得按"并行度 × 每并行度行数"的规则推算预期总行数,避免误判。

常见问题排查

  • connector loading error/ 找不到插件:回到 Step 2 确认connectors/下确实存在两个 JAR,且config/plugin_config中插件名拼写无误(注意是connector-fake/connector-console,不是fake/console)。
  • 下载失败或校验和不匹配:检查网络能否访问 Maven 中央仓库;如需走内网镜像,设置SEATUNNEL_MAVEN_REPOSITORY后重新执行install-plugin.sh
  • 未打印output rowType:多半是配置解析阶段就出错,可先用-d static做一次纯配置校验,观察是否报 schema 或插件配置相关错误。
  • 输出行数不符:核对parallelismrow.num的乘积关系,以及 transform/sink 的plugin_input通道是否与上游plugin_output一致。

下一步

本示例成功后,说明你的 SeaTunnel 本地基础路径已完全打通,可以进入真实管道开发:

  • 完整的本地引擎走查,继续阅读 Quick Start With SeaTunnel Engine(默认引擎,通常是最短的成功路径);若使用 Flink 或 Spark 作为执行引擎,可分别参考 Quick Start With Flink 与 Quick Start With Spark;
  • 首个经过验证的 Source→Sink 实战案例,从 MySQL CDC to Kafka 开始;
  • 其他管道形态可继续探索:
    • MySQL CDC to Doris
    • JDBC to S3
    • Kafka to Iceberg
    • Http to JDBC
    • File to StarRocks
    • Multi-table CDC

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

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

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

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

立即咨询