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.sh、bin/seatunnel.cmd、config/、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-fake与connector-console两项:
--seatunnel-connectors-- connector-fake connector-console --end--关于plugin_config需要说明两点(均可在仓库中直接核对):
- 仓库中真实的 config/plugin_config 使用
--connectors-v2--作为段标记,并包含了全部已注册连接器的 artifactId(如connector-jdbc、connector-kafka、connector-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 直连下载(需要
curl、mktemp以及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能力; - 对于
SNAPSHOT、LATEST、RELEASE及版本区间等无法直接解析的动态版本,脚本会自动切换为 Maven 方式下载。
若脚本运行后connectors/目录下能看到connector-fake-3.0.0.jar与connector-console-3.0.0.jar,则插件安装成功。Windows 用户请使用bin\install-plugin.cmd(该脚本使用捆绑的 Maven Wrapper,无需单独安装 Maven)。
Step 3:编写一个最小任务配置
将下面的配置保存为config/v2.batch.config.template或任意本地文件。它是标准的 SeaTunnel HOCON 配置,由env、source、transform、sink四大块组成:
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.num | 5 | 每个并行度生成的数据总行数(示例显式设为 16) |
split.num | 1 | 每个并行度由 enumerator 切分的 split 数量 |
split.read-interval | 1 | reader 两次 split 读取之间的间隔(毫秒) |
string.length | 5 | 生成的 string 类型字段长度 |
map.size/array.size/bytes.length | 5 | 对应复杂类型的生成尺寸 |
string.template | 无 | 若配置,则 string 字段从模板列表中随机选取 |
int.min/int.max | 0/Integer.MAX_VALUE | 整型生成范围,其他数值类型同理(tinyint、smallint、bigint、float、double等均有min/max对) |
rows | 无 | 显式指定要输出的行列表(每行含kind与fields),优先级高于随机生成 |
string.fake.mode等 | RANGE | 生成模式,可选RANGE(区间随机)或TEMPLATE(模板选取) |
auto.increment.enabled | false | 是否启用自增 ID 生成,配合auto.increment.start(默认 1)使用 |
row.num的"每并行度"语义很重要:若你设置parallelism = 2且row.num = 16,实际会生成 32 行数据,因为 FakeSourceOptions.java 中明确描述其为 "The total number of data generated per degree of parallelism"。此外schema.fields定义了输出表的字段结构(此处为name: string与age: 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)的取值被限定为local与cluster两种:-m local表示由客户端本地拉起引擎执行,无需预先部署 SeaTunnel Engine 服务集群,这也是新用户最推荐的入门方式;cluster模式则要求先部署引擎服务。-e/--deploy-mode自 2.3.1 起已被标记为 deprecated,建议统一使用-m/--master。
除--config与-m外,该参数类还提供了若干实用选项,可用于后续调试:
-d/--dry-run:在不真正运行 Sink 的前提下做校验或预览,支持static(仅静态校验配置)、connect(校验连接)、sample(采样预览数据)三种模式;--sample-limit:sample模式下每个 Source 最多转发行数(默认 10,上限 10000);-l/--list:列出作业状态;-j/--job-id:按 JobId 查询作业状态。
预期验证结果
任务提交后,应按以下四条标准逐项核对(全部满足即说明本机基础链路健康):
- 进程正常启动,无连接器加载错误:说明
connector-fake与connector-console两个 JAR 被正确发现并加载; - 控制台打印
output rowType行:这一行由 ConsoleSinkWriter.java 在初始化时以log.info("output rowType: {}", ...)输出,内容应展示经过FieldMapper映射后的字段,即age与new_name,这是验证 transform 生效的直接证据; - 控制台打印来自
ConsoleSinkWriter的 16 行数据:每条记录形如subtaskIndex=0 rowIndex=N: SeaTunnelRow#tableId=... SeaTunnelRow#kind=INSERT : <字段值>(见 ConsoleSinkWriter.java),这里出现 16 行正是因为parallelism = 1且row.num = 16; - 批任务在写完所有行后正常退出:
BATCH模式下数据有界,作业完成后客户端应正常结束并返回。
如果你修改了parallelism或row.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 或插件配置相关错误。 - 输出行数不符:核对
parallelism与row.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),仅供参考