SeaTunnel 基于 Spark 引擎的快速开始:部署、作业配置与源码级原理解析

发布时间:2026/9/20 5:01:34
SeaTunnel 基于 Spark 引擎的快速开始:部署、作业配置与源码级原理解析
SeaTunnel 基于 Spark 引擎的快速开始部署、作业配置与源码级原理解析【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本篇技术指南以 Apache SeaTunnel 官方文档《Spark 引擎快速开始》为骨架面向已经明确要将 SeaTunnel 跑在 Spark 上的团队系统讲解从部署 SeaTunnel、配置 Spark 环境到编写 HOCON 作业配置文件、提交运行的完整闭环并深入 Spark starter 与转换层源码解释spark-submit命令如何被生成、spark.前缀参数如何生效以及 SeaTunnel 作业如何被适配到 Spark 的 datasource 执行模型。读完本文你将能够独立完成一次基于 Spark 的 SeaTunnel 端到端作业从 FakeSource 到 Console的搭建、运行与验证并为后续接入真实连接器奠定基础。什么时候选择 Spark 作为执行引擎SeaTunnel 支持多种执行引擎Spark 是其中一种成熟选项。根据 SeaTunnel 运行在 Spark 上 的说明以下场景通常更适合使用 Spark团队已经在生产中稳定运行 Spark 集群周边任务主要以批处理为主希望 SeaTunnel 与既有 Spark 生态和部署方式保持一致。需要注意使用 Spark 或 Flink 运行 SeaTunnel 任务时无需部署 SeaTunnel 引擎Zeta服务集群这与内置 Zeta 引擎的部署方式不同参见 部署。如果你是第一次评估 SeaTunnel、且没有必须使用 Spark 的前提建议先从 SeaTunnel 引擎快速开始 入手那是默认引擎下最短的验证路径。步骤 1部署 SeaTunnel 及连接器在开始前请确保已经按照 部署 中的描述下载并部署了 SeaTunnel。部署前的软硬件前提包括安装 Java 8 或 11其他高于 Java 8 的版本理论上也可以工作并正确设置JAVA_HOME。下载发行包后由于从 2.2.0-beta 版本开始二进制包不再默认携带连接器依赖首次使用需要安装连接器插件sh bin/install-plugin.shSeaTunnel 通过config/plugin_config文件控制要安装哪些插件。为了让本文示例链路FakeSource→FieldMapper→Console正常工作至少需要connector-fake和connector-console两个插件plugin_config内容如下--seatunnel-connectors-- connector-fake connector-console --end--你可以在${SEATUNNEL_HOME}/connectors/plugins-mapping.properties下查看所有支持的连接器及其对应的 plugin_config 配置名称也可以从 Maven 仓库手动下载连接器 JAR放入${SEATUNNEL_HOME}/connectors/目录2.3.5 之前为connectors/seatunnel目录。步骤 2部署并配置 SparkSeaTunnel 的 Spark starter 需要基于真实的 Spark 发行版运行因此第二步是准备 Spark 环境。下载 Spark请先下载 Spark 发行包需要版本 2.4.0。更多信息可参考 Spark 官方文档的 Standalone 模式安装指南。配置 SeaTunnel 中的 SPARK_HOME修改config/seatunnel-env.sh中的设置该文件基于引擎的安装路径将SPARK_HOME修改为 Spark 的部署目录即可。仓库内置的 config/seatunnel-env.sh 默认值如下# Home directory of spark distribution. SPARK_HOME${SPARK_HOME:-/opt/spark}也就是说如果 Spark 安装在其他目录你需要把这里改成实际路径例如SPARK_HOME/usr/local/spark。这一环境变量是后续生成spark-submit命令的关键前提。步骤 3添加作业配置文件来定义作业SeaTunnel 的大多数作业通过声明式配置完成你通常不需要先写代码。编辑config/v2.streaming.conf.template本文示例所用它决定了 SeaTunnel 启动后数据输入、处理和输出的方式及逻辑。下面是完整示例配置与本文要运行的示例应用程序一致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 } }四个配置块的角色参照 作业配置指南一个标准作业由四个顶层块组成职责划分如下env控制作业如何执行。通用参数包括job.modeBATCH或STREAMING、parallelism作业默认并行度、job.name可选作业显示名称、checkpoint.interval流作业或 exactly-once 场景下的 checkpoint 间隔。Flink/Spark 的引擎专属参数也放在这里。source定义数据从哪里来包括连接器名称、连接参数、读取范围表、topic、路径或查询、schema/format 相关参数以及用于下游显式引用的plugin_output。transform可选的中间处理层支持字段重命名/映射、行过滤、RowKind 处理、SQL 转换、写入前校验等。若链路不需要中间转换可以整体省略。sink定义数据最终写到哪里包括连接器名称、连接参数、目标表/topic/路径、写入语义以及声明消费哪个上游输出的plugin_input。plugin_input 与 plugin_output 约定plugin_output用于给 source 或 transform 的输出命名plugin_input用于让 transform 或 sink 指向某个上游输出。在存在多个 source、一个 transform 输出写入多个 sink、或链路较复杂需要保证可读性时这两个字段尤为重要。链路只有一个上游时 SeaTunnel 往往可以依赖默认约定流转但从可维护性角度仍建议显式命名。Spark 专属配置的写法如果使用 Spark 引擎Spark 专属作业参数同样写在env块中并使用spark.前缀参见 SeaTunnel 运行在 Spark 上env { spark.app.name example spark.sql.catalogImplementation hive spark.executor.memory 2g spark.executor.instances 2 spark.yarn.priority 100 spark.dynamicAllocation.enabled false }这些env中的键值对会被 Spark starter 解析为--conf keyvalue形式的参数透传给spark-submit源码依据见下文 SparkStarter 如何生成 spark-submit 命令。关于配置的基本概念可继续阅读配置的基本概念数据集编排与 Transform 参数可参考 作业配置指南 与 Transform 通用参数。步骤 4运行 SeaTunnel 应用程序SeaTunnel 针对不同的 Spark 大版本提供了独立的启动脚本。启动命令如下Spark 2.4.xcd apache-seatunnel-${version} ./bin/start-seatunnel-spark-2-connector-v2.sh \ --master local[4] \ --deploy-mode client \ --config ./config/v2.streaming.conf.templateSpark 3.x.xcd apache-seatunnel-${version} ./bin/start-seatunnel-spark-3-connector-v2.sh \ --master local[4] \ --deploy-mode client \ --config ./config/v2.streaming.conf.template命令行参数解析启动脚本最终会把上述参数交给SeaTunnelSpark主类处理。仓库中的 SparkCommandArgs.java 定义了这些参数的语义参数含义默认值-m/--masterSpark master支持spark://host:port、mesos://host:port、yarn、k8s://https://host:port、locallocal[*]-e/--deploy-modeSpark 部署模式仅支持cluster、clientclient--config作业配置文件路径无其他通用参数还包括--check仅校验配置不执行、--encrypt/--decrypt配置加密/解密以及-i变量注入等。Spark on YARN 的典型用法参见 SeaTunnel 运行在 Spark 上# Spark on YARN 集群模式 ./bin/start-seatunnel-spark-3-connector-v2.sh --master yarn --deploy-mode cluster --config config/example.conf # Spark on YARN 客户端模式 ./bin/start-seatunnel-spark-3-connector-v2.sh --master yarn --deploy-mode client --config config/example.confSparkStarter 如何生成 spark-submit 命令从源码看SparkStarter.java 并不是直接执行作业而是“根据 SeaTunnel 配置生成一条spark-submit命令”。其buildFinal()方法约 L170-L199依次组装${SPARK_HOME}/bin/spark-submit可执行文件路径--class org.apache.seatunnel.core.starter.spark.SeaTunnelSpark主类--name、--master、--deploy-mode等命令参数--jarslib 目录 JAR 按作业配置发现的连接器 JAR 第三方 JAR将env块解析出的 Spark 配置逐条追加为--conf keyvaluegetSparkConf()方法L132-L138应用 JAR 与--config配置文件路径。同时getInstance()L90-L106会根据--deploy-mode选择不同的 Starter 子类ClientModeSparkStarter会把 driver 相关内存、Java 选项等映射为--driver-memory等专用参数ClusterModeSparkStarter则会把插件目录打成 tar.gz 并通过--files随作业分发到集群L327-L334。理解这条链路有助于在排查启动问题时定位参数来源。查看输出当运行命令时你可以在控制台中看到 SeaTunnel 的输出这可以视为命令运行成功或失败的标志。SeaTunnel 控制台将会打印类似如下的日志fields : name, age types : STRING, INT row1 : elWaB, 1984352560 row2 : uAtnp, 762961563 row3 : TQEIB, 2042675010 row4 : DcFjo, 593971283 row5 : SenEb, 2099913608 row6 : DHjkg, 1928005856 row7 : eScCM, 526029657 row8 : sgOeE, 600878991 row9 : gwdvw, 1951126920 row10 : nSiKE, 488708928 row11 : xubpl, 1420202810 row12 : rHZqb, 331185742 row13 : rciGD, 1112878259 row14 : qLhdI, 1457046294 row15 : ZTkRx, 1240668386 row16 : SGZCr, 94186144从日志可以看到两处关键信息fields : name, age对应FakeSource中定义的 schema输出字段变为new_name与age正是FieldMapper中name new_name映射生效的结果。16 行数据与row.num 16一致。源码级原理SeaTunnel 到 Spark 的转换层如果你想知道 SeaTunnel 的 connector API 是如何在 Spark 上运行的Spark 转换层 文档给出了完整解释。SeaTunnel 并不是简单地把 Spark 当作“另一个 Flink”而是把 SeaTunnel 契约重新解释成 Spark 可接受的执行模型。从概念上看Spark 路径的高层映射关系如下SeaTunnelSource - Spark source adapter - Spark datasource runtime SeaTunnelSink - Spark sink adapter - Spark datasource writer runtime SeaTunnel schema/types - Spark schema/types - InternalRow execution其中最重要的三类问题是Source 侧适配以 Spark 期望的形式暴露 schema、从 SeaTunnel split 规划 Spark partitions、为每个 partition 创建 reader、把 SeaTunnel 输出转换成 SparkInternalRow。与 Flink 相比Spark 更强调“规划完成的 partitions reader 执行”而非长期存在的 enumerator/coordinator 模型。Sink 侧适配创建 writer factory、在 executor 与 driver 之间传递 commit message、协调 commit 与 abort 路径、把 SeaTunnel 的 retry 语义映射到 Spark 可接受的行为。当 sink 不是 append-only、依赖幂等或事务提交时尤其关键。Schema 与 Row 转换把CatalogTable/TableSchema、SeaTunnelDataType、SeaTunnelRow映射为 Spark 侧的StructType、Spark SQL data type 和InternalRow。decimal、timestamp、嵌套类型和 nullability 是最容易出问题的边界。值得注意的还有版本拆分Spark 2.4 与 Spark 3.x 的 datasource API 并不完全一致因此仓库在seatunnel-translation/seatunnel-translation-spark/下保留了seatunnel-translation-spark-2.4与seatunnel-translation-spark-3.3等独立适配模块这就是为什么快速开始要区分start-seatunnel-spark-2-connector-v2.sh与start-seatunnel-spark-3-connector-v2.sh两个启动脚本。可优先关注SparkSink、SparkDataSourceWriter、SeaTunnelInputPartitionReader等实现类。如果 schema 转换不一致、InternalRow转换出错、datasource writer commit 行为异常问题往往不在 connector 而在转换层内部排查时应注意这一层。从源码仓库运行示例如果你是直接在源码仓库中运行示例对应的模块与入口类记载于 SeaTunnel 运行在 Spark 上为模块seatunnel-examples/seatunnel-spark-connector-v2-example示例入口类org.apache.seatunnel.example.spark.v2.SeaTunnelApiExample从示例迁移到真实作业跑通 FakeSource → Console 示例后推荐按下面顺序逐步替换为真实链路参见 作业配置指南保留示例中的env基本结构用真实 Source 替换FakeSource用真实 Sink 替换Console只有在源表结构和目标结构不直接匹配时再补充 transform按连接器要求补充驱动 jar 或额外依赖例如 JDBC 连接器需要把数据库驱动放入${SEATUNNEL_HOME}/lib/。运行前建议按以下清单检查Java 与JAVA_HOME已正确配置所需插件已经安装plugin_configinstall-plugin.sh第三方驱动 jar 已就位Source 凭据与网络访问正确目标表、topic 或路径已提前创建job.mode与连接器能力相匹配SPARK_HOME指向正确的 Spark 安装目录。进一步阅读开始编写你自己的配置文件选择想要使用的连接器并根据连接器的文档配置参数想了解更多 SeaTunnel 运行在 Spark 上的信息请参阅基于 Spark 的 SeaTunnel想理解 SeaTunnel API 如何被适配到 Spark请继续阅读 Spark 转换层想了解引擎差异与选型参见引擎概览如果你的场景不强制要求 SparkSeaTunnel 内置的Zeta引擎默认推荐提供了更短的本地验证路径可回到 SeaTunnel 引擎快速开始。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考