ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

SeaTunnel 本地跑通第一个任务:FakeSource + FieldMapper + Console 最小链路实战指南

SeaTunnel 本地跑通第一个任务:FakeSource + FieldMapper + Console 最小链路实战指南 数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载导读本文聚焦一个目标——用最短路径在本地把 SeaTunnel 真正跑起来验证「安装 → 插件下载 → 配置解析 → 执行引擎」整条链路是否正常。全程不依赖 MySQL、Kafka 或对象存储任何具备 Java 环境的机器都能复现。读完本文你将掌握最小可运行配置的写法、install-plugin.sh的插件裁剪技巧、本地模式-m local的启动方式以及如何从控制台日志验证批任务执行结果。一、这篇指南要解决什么问题SeaTunnel 是一个多模态、高性能、分布式的海量数据集成工具连接器种类多达上百个但第一次上手时并不需要它们全部。run-your-first-job这篇指南即本文的原始依据位于 docs/zh/getting-started/locally/run-your-first-job.md的核心价值在于用两个最轻量的连接器connector-fake与connector-console搭建一条纯本地的数据流水线一次性确认以下四件事都正常本地部署是否成功bin/seatunnel.sh存在且可执行插件安装机制是否工作install-plugin.sh按需下载连接器 JAR配置文件HOCON 格式能否被正确解析内置执行引擎能否完成「读取 → 转换 → 写出」并正常退出。链路走通之后再切换到真实数据源MySQL、Kafka、对象存储等和真实目标端就有了可靠的基座。二、步骤 1先完成本地部署运行第一个任务的前提是完成本地部署。请先阅读 本地部署指南其中明确了环境要求安装 Java 8 或 11其他高于 Java 8 的版本理论上也可以工作并设置JAVA_HOME获取发行包从 Apache 下载页面获取seatunnel-version-bin.tar.gz解压后得到形如apache-seatunnel-version/的目录。完成部署后先确认 SeaTunnel 目录下存在bin/seatunnel.shls bin/seatunnel.shbin/seatunnel.sh是 SeaTunnel 任务的统一启动入口后续所有本地运行命令都经由它提交。如果这一步就找不到脚本说明解压不完整或目录结构不对需要回到部署步骤检查。三、步骤 2只安装这篇示例真正需要的插件3.1 为什么需要安装插件从 2.2.0-beta 版本开始SeaTunnel 的二进制发行包不再默认携带连接器依赖。也就是说第一次使用时必须手动安装连接器否则启动任务时会因为找不到连接器类而报错。安装动作由发行包根目录下的bin/install-plugin.sh完成。开发者说明如果你需要验证未发布代码、调试 SeaTunnel 源码或构建自定义发行包请参考 搭建开发环境 末尾指向的开发者文档docs/zh/developer/setup.md。本地部署指南默认面向使用官方二进制发行包的用户。3.2 收敛 plugin_config仓库自带的 config/plugin_config 列出了全部可用连接器connector-amazondocumentdb、connector-cdc-mysql、connector-jdbc、connector-kafka……数百行但绝大多数对首个示例无用。请把它收敛成只有下面两个插件--seatunnel-connectors-- connector-fake connector-console --end--connector-fake示例数据源不连接任何外部系统按配置批量生成模拟数据connector-console示例目标端把收到的每一行数据打印到控制台。说明config/plugin_config中的--connectors-v2--与文档示例中的--seatunnel-connectors--都是区块分隔标识核心是只保留你真正需要的连接器名称行。你可以在${SEATUNNEL_HOME}/connectors/plugins-mapping.properties对应仓库根目录的 plugin-mapping.properties中查看所有支持连接器的 plugin_config 配置名称。3.3 执行安装命令并验证cd ${SEATUNNEL_HOME} sh bin/install-plugin.sh ls connectors | rg connector-(fake|console)install-plugin.sh会读取plugin_config仅下载列表内的连接器 JAR 到${SEATUNNEL_HOME}/connectors对于正式发布的连接器版本脚本通过 HTTPS 直接下载 JAR 及其校验文件因此 Linux 和 macOS 不需要 Maven需要curl、mktemp以及sha512sum、sha1sum、shasum或openssl中的任意一个Windows 的install-plugin.cmd仍使用发行包内置的 Maven Wrapper需要指定版本时执行sh bin/install-plugin.sh version例如sh bin/install-plugin.sh 3.0.0如需 Maven 兼容的 HTTPS 镜像可通过SEATUNNEL_MAVEN_REPOSITORY环境变量指定仓库根地址。最后一条ls connectors | rg connector-(fake|console)应当同时输出connector-fake和connector-console两个目录证明插件已就位。四、步骤 3编写最小可运行配置把下面的配置保存为config/v2.batch.config.template覆盖仓库自带的模板 config/v2.batch.config.template或者保存为你自己的本地配置文件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 } }这个配置由四个区块构成正好对应 SeaTunnel 的数据处理模型区块作用本示例的关键配置env全局运行环境parallelism 1控制并行度job.mode BATCH声明批模式source数据源FakeSource生成 16 行数据schema 定义namestring、ageint两个字段transform数据转换FieldMapper做字段映射把name改名为new_namesink数据目标Console把结果打印到控制台区块之间通过plugin_output/plugin_input形成数据管道fakeFakeSource 输出→fake1FieldMapper 输出→Console消费fake1。4.1 FakeSource 参数背后源码如何解释这些配置row.num、schema等参数并非魔法字符串它们在 FakeSourceOptions.java 中被显式声明为 Option 定义并在 FakeConfig.java 的buildWithConfig中完成解析row.num每个并行度生成的数据总行数默认值 5本示例设为 16因此在parallelism 1下总共生成 16 行split.num每个并行度由枚举器生成的 split 数量默认 1split.read-interval单个 reader 两次 split 读取之间的间隔毫秒默认 1map.size/array.size/bytes.length/string.length生成 map、array、bytes、string 类型数据时的尺寸/长度默认均为 5各类数值字段tinyint/smallint/int/bigint/float/double均支持xxx.min、xxx.max控制随机范围或xxx.template从给定列表随机取值rows直接指定要逐行产出的数据列表含kind与fields比row.num更精确。schema.fields则通过CatalogTableUtil.buildWithConfig构建出表的 Catalog 结构定义了 FakeSource 产出的行类型namestring、ageint。4.2 FieldMapper 转换字段重命名的实现原理FieldMapper的配置键field_mapper在 FieldMapperTransformConfig.java 中被定义为OptionMapString, String语义为「输入字段 → 输出字段」的映射关系。本示例中field_mapper { age age # age 保持不变 name new_name # name 重命名为 new_name }在 FieldMapperTransform.java 的构造器中会逐个校验field_mapper的 key 是否存在于输入行类型中若找不到输入字段会抛出cannotFindInputFieldsError异常这解释了为什么配置写错字段名会直接启动失败而不是运行期报错。执行转换时transformRow方法它按映射顺序从输入行取字段值组装出新的SeaTunnelRow并保留原始行的rowKind、tableId等元信息——这些细节决定了下游Console输出的日志格式。五、步骤 4用本地模式运行进入发行包目录用-m local指定本地模式运行cd apache-seatunnel-${version} ./bin/seatunnel.sh --config ./config/v2.batch.config.template -m local关键点说明-m local指定运行模式为local本地模式任务在提交的进程内直接执行不需要预先部署 SeaTunnel 引擎Zeta集群如果想跑在独立 Zeta 集群上需要先部署引擎服务参考 SeaTunnel 引擎快速开始--config指向步骤 3 保存的配置文件路径相对于当前目录也可以使用 Flink / Spark 引擎运行但那样需要额外的运行时环境且并非本示例的最短路径。六、验证结果任务跑完后对照以下四点逐项检查全部满足即说明本地基础链路正常任务可以正常启动没有 connector 加载错误——即插件安装正确、类加载完整控制台打印映射后字段的output rowType行——这一行来自ConsoleSinkWriter构造时的日志log.info(output rowType: {}, ...)见 ConsoleSinkWriter.java它会按字段名类型的格式输出整个行结构例如new_namestring, ageint正好证明FieldMapper的字段重命名已生效控制台打印 16 行ConsoleSinkWriter输出——ConsoleSinkWriter.write方法对每一行数据输出subtaskIndex{} rowIndex{} SeaTunnelRow#tableId{} SeaTunnelRow#kind{} : 字段值列表row.num 16对应恰好 16 条批任务在写完全部数据后正常退出——job.mode BATCH下任务有界数据处理完毕即结束进程不会像流任务那样常驻。如果以上四点都通过说明「插件安装 → 配置解析 → 引擎调度 → 数据流转」整条本地链路是健康的可以放心切换到真实数据源和真实目标端。七、调试排查建议连接器加载报错优先回到步骤 2确认ls connectors | rg connector-(fake|console)有输出且plugin_config中没有写错连接器名称output rowType与预期字段不一致检查transform.field_mapper的 key 是否都是输入 schema 中真实存在的字段名否则FieldMapperTransform会在初始化期抛异常行数与预期不符row.num是每并行度的行数总行数 row.num × parallelism调整env.parallelism会成倍改变输出行数。八、下一步从最小链路走向真实场景本地最小链路跑通后按你的数据链路形态选择后续路线想看完整本地引擎链路继续阅读 SeaTunnel 引擎快速开始想看一条刚完成端到端核对的新教程先看 MySQL CDC 到 Kafka其他典型链路示例均在docs/zh/getting-started/recipes/下MySQL CDC 到 DorisJDBC 到 S3Kafka 到 IcebergHttp 到 JDBCFile 到 StarRocks多表 CDC切换到真实场景时只需把source/sink区块替换为对应连接器的配置例如 MySQL CDC Source、Kafka Sinkenv、transform、plugin_input/plugin_output的数据流骨架保持不变——这正是先用 FakeSource Console 验证链路的长期价值所在。赞分享数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载相关推荐SeaTunnel 首个任务实战基于 FakeSource 与 Console 的本地全链路验证指南SeaTunnel 首个任务实战基于 FakeSource 与 Console 的本地全链路验证指南 本篇指南带你走完 Apache SeaTunnel 的数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel 本地快速上手指南最短路径跑通你的第一个数据同步任务SeaTunnel 本地快速上手指南最短路径跑通你的第一个数据同步任务 本篇指南围绕 Apache SeaTunnel 的 Local Quick Start数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel 本地快速开始从零部署到跑通第一个数据同步任务的完整指南SeaTunnel 本地快速开始从零部署到跑通第一个数据同步任务的完整指南 SeaTunnel 是一个多模态、高性能、分布式的海量数据集成工具通过声明式配置数据集成ETL大数据批处理流处理变更数据捕获创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表