
SeaTunnel 在 Spark 上的快速入门部署、配置与运行你的第一个同步作业【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本指南面向已经拥有或计划使用Apache Spark 集群、希望让 SeaTunnel 作业无缝融入既有批处理 / 混合负载环境的团队完整讲解如何在 Spark 引擎上部署 SeaTunnel、配置环境、编写作业定义文件并用官方启动脚本提交第一个数据同步任务。读完本文你将掌握从 Spark 版本选型、SPARK_HOME配置、HOCON 作业配置编写到spark-submit底层命令生成机制的完整实战链路。适用前提若你是第一次评估 SeaTunnel 且没有必须使用 Spark 的诉求建议先走默认引擎 Zeta 的 Quick Start With SeaTunnel Engine那是本地验证的最短路径。使用 Spark 运行同步任务时无需部署 SeaTunnel EngineZeta服务集群这是 Deployment 文档 中明确说明的差异点。一、开跑前的准备阅读本文假定你已对 SeaTunnel 的引擎选择与作业配置有一定背景认知。在继续之前建议按需阅读以下文档它们与本文构成完整的 Spark 运行链路Engine Overview了解 SeaTunnel 支持的多引擎架构与各自适用场景SeaTunnel With Spark理解何时选择 Spark以及 Spark 专属配置项Job Configuration Guide掌握env/source/transform/sink四段式作业结构的通用规则。二、Step 1部署 SeaTunnel 与连接器插件运行 Spark 模式前需要先完成 SeaTunnel 本体的下载与部署详见 Deployment核心要点如下。1. 环境准备安装 Java8 或 11高于 8 的版本理论上也可用并正确设置JAVA_HOME。2. 下载二进制发行包从 SeaTunnel 下载页获取seatunnel-version-bin.tar.gz或通过终端下载解压export version3.0.0 wget https://archive.apache.org/dist/seatunnel/${version}/apache-seatunnel-${version}-bin.tar.gz tar -xzvf apache-seatunnel-${version}-bin.tar.gzWindows 用户可下载.zip包后用文件管理器或 PowerShell 解压。3. 安装连接器插件自 2.2.0-beta 起二进制包默认不再内置连接器依赖首次使用需执行插件安装脚本sh bin/install-plugin.sh关键细节不需要安装全部连接器。通过修改config/plugin_config即可指定所需插件。以本文示例作业为例只需connector-fake与connector-console两个插件配置如下--seatunnel-connectors-- connector-fake connector-console --end--全部受支持连接器及其对应的plugin_config名称可在${SEATUNNEL_HOME}/connectors/plugins-mapping.properties中查到仓库根目录的 plugin-mapping.properties 即该类映射的源码级样例。也可以手动从 Maven 仓库下载连接器 JAR放到${SEATUNNEL_HOME}/connectors/目录2.3.5 之前版本需放在connectors/seatunnel目录。安装脚本支持指定版本如sh bin/install-plugin.sh 3.0.0并通过 HTTPS 直接下载 JAR 与校验和需curl、mktemp以及sha512sum/sha1sum/shasum/openssl之一也可通过环境变量SEATUNNEL_MAVEN_REPOSITORY指向 HTTPS Maven 镜像或用SEATUNNEL_PLUGIN_DOWNLOAD_METHODmaven保留 Mavensettings.xml的镜像、认证仓库、代理等行为。三、Step 2部署并配置 Spark1. 下载 Spark首先下载 Apache Spark 发行包要求版本 2.4.0。安装方式可参考 Spark 官方 Standalone 部署文档。SeaTunnel 对 Spark 主版本线是分别适配的对应两套启动脚本详见 Step 4这也是仓库中 Spark 翻译层按 2.4 与 3.x 分模块维护的原因。2. 配置 SeaTunnel 的环境脚本编辑${SEATUNNEL_HOME}/config/seatunnel-env.sh将SPARK_HOME指向 Spark 部署目录。仓库中的 config/seatunnel-env.sh 给出了默认值与可覆盖方式# Home directory of spark distribution. SPARK_HOME${SPARK_HOME:-/opt/spark}即在未显式导出SPARK_HOME时默认取/opt/spark请按实际部署路径覆盖该变量。该脚本还包含FLINK_HOMEFlink 模式使用以及 Metalake 相关的METALAKE_ENABLED/METALAKE_TYPE/METALAKE_URL配置与本主题无关的部分无需改动。补充两个 Spark 启动脚本在运行时都会 source 该环境文件——见 start-seatunnel-spark-2-connector-v2.sh 与 start-seatunnel-spark-3-connector-v2.sh 中的if [ -f ${CONF_DIR}/seatunnel-env.sh ]; then . ${CONF_DIR}/seatunnel-env.sh; fi因此SPARK_HOME的配置必须写入该文件或提前导出为环境变量。四、Step 3编写作业配置文件config/v2.streaming.conf.template是 SeaTunnel 启动后定义数据输入、处理与输出逻辑的作业文件。仓库中的 config/v2.streaming.conf.template 是流式示例而本文示例作业与官方 Spark 快速入门一致使用批处理模式完整内容如下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控制作业的执行方式parallelism 1指定默认并行度job.mode BATCH声明批处理模式source定义数据来源FakeSource生成 16 行模拟数据字段nameSTRING与ageINTtransform对在途数据做变换可选FieldMapper重命名字段将name映射为new_namesink定义数据去向Console将结果打印到控制台关键参数说明plugin_output/plugin_input这是理解 SeaTunnel 数据流的核心约定。plugin_output为 source/transform 产出的数据流命名plugin_input告诉下游 transform/sink 消费哪条上游流。本示例中FakeSource输出名为fakeFieldMapper消费fake并输出fake1Console消费fake1。当作业只有一个上游路径时SeaTunnel 通常能按默认约定自动串联但显式命名更利于多源、多分支作业的可读性详见 Job Configuration Guide。row.numFakeSource 生成的行数定义于 FakeSourceOptions.java本示例为 16。schema.fields声明字段名与类型的映射name为string、age为int对应输出日志中的types : STRING, INT。field_mapperFieldMapper 的字段重命名映射age age保持不变name new_name将原字段name重命名为new_name。因此最终控制台输出仍只有两列重命名后的new_name与age。job.modeBATCH或STREAMING。注意仓库自带的v2.streaming.conf.template是流式示例job.mode STREAMING、checkpoint.interval 2000与本文批处理示例在env块上不同运行前请按需调整。更完整的配置概念请查阅 Config Concept变换类插件的通用参数见 Transform Common Options。五、Step 4运行 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命令行参数含义参数含义本示例取值--masterSpark 集群 master 地址local[4]表示本地 4 核运行生产环境可替换为yarn、spark://host:port等--deploy-mode部署模式client驱动在提交端运行集群模式用cluster--config作业配置文件路径./config/v2.streaming.conf.template启动脚本的底层机制动态生成 spark-submit 命令两个启动脚本并非直接调用 Spark API 提交作业而是先以org.apache.seatunnel.core.starter.spark.SparkStarter为入口执行参数解析在 stdout 输出一条拼装好的spark-submit命令字符串再由脚本eval执行。以 start-seatunnel-spark-3-connector-v2.sh 为例CMD$(java ${JAVA_OPTS} -cp ${CLASS_PATH} ${APP_MAIN} ${args}) EXIT_CODE$? || EXIT_CODE$? ... elif [ ${EXIT_CODE} -eq 0 ]; then echo Execute SeaTunnel Spark Job: $(echo ${CMD} | tail -n 1) eval $(echo ${CMD} | tail -n 1)其中APP_MAINorg.apache.seatunnel.core.starter.spark.SparkStarterAPP_JAR${APP_DIR}/starter/seatunnel-spark-3-starter.jar。脚本还会自动加载config/seatunnel-env.sh并设置 log4j2 配置文件与日志目录。在 SparkStarter.java 的buildFinal()方法中可以看到生成的命令骨架${SPARK_HOME}/bin/spark-submit --class org.apache.seatunnel.core.starter.spark.SeaTunnelSpark --name job-name --master master --deploy-mode client|cluster --jars 连接器插件 JAR 列表 --files 附属文件 --conf spark 配置项 starter jar 路径 --config 作业配置文件这解释了为什么必须正确配置SPARK_HOME脚本最终提交的spark-submit命令就是${SPARK_HOME}/bin/spark-submit。同时getPluginIdentifiers(...)会根据作业配置中的 source/transform/sink 插件名通过SeaTunnelSourcePluginDiscovery/SeaTunnelSinkPluginDiscovery自动收集对应连接器 JAR 及其依赖并通过--jars注入提交命令这正是 Step 1 必须安装好对应连接器插件的直接原因。六、查看输出验证作业是否成功运行命令后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与types : STRING, INT来自作业schema定义与schema.fields完全对应共16行记录row1至row16与row.num 16一致每行由两个字段构成经FieldMapper重命名后输出列语义不变字段名变为new_name与age随机内容来自 FakeSource 的模拟生成逻辑。七、深入SeaTunnel 如何适配 Spark 运行时若你想理解 SeaTunnel Connector API 如何被翻译到 Spark 执行模型可阅读 Spark Translation Layer。核心思路是把 SeaTunnel 的契约重新解释为 Spark 原生概念而不是简单做接口改名。高层的映射关系为SeaTunnelSource - Spark source adapter - Spark datasource runtime SeaTunnelSink - Spark sink adapter - Spark datasource writer runtime SeaTunnel schema/types - Spark schema/types - InternalRow execution源码结构上翻译层按 Spark 主版本分模块实现与两套启动脚本一一对应seatunnel-translation/seatunnel-translation-spark/seatunnel-translation-spark-common/公共部分如InternalRowConverter、SeaTunnelRowConverter、TypeConverterUtils等行列与类型转换工具seatunnel-translation/seatunnel-translation-spark/seatunnel-translation-spark-2.4/Spark 2.4 适配如SparkSink、SparkDataSourceWriter、SeaTunnelInputPartitionReader等seatunnel-translation/seatunnel-translation-spark/seatunnel-translation-spark-3.3/Spark 3.3 适配如SeaTunnelSparkSource、SeaTunnelBatch/SeaTunnelMicroBatch、SeaTunnelSparkSink、SeaTunnelWrite等。翻译层最敏感的边界包括source 侧的分区规划从 SeaTunnel split 信息映射为 Spark 输入分区、schema 与行的转换CatalogTable/TableSchema、SeaTunnelDataType、SeaTunnelRow到StructType、Spark SQL 类型、InternalRow尤其关注 decimal、timestamp、嵌套类型与空值语义以及 sink 侧的提交 / 中止路径writer commit message、幂等与事务语义的桥接。常见故障往往表现为 schema 转换不匹配、InternalRow转换异常或 writer 提交行为差异且根因可能藏在翻译层而非连接器本身。Spark 专属配置spark.前缀当引擎为 Spark 时Spark 专属的作业参数在env块中使用spark.前缀例如 spark.md 中的示例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 }这些键值最终由SparkStarter.appendSparkConf(...)转换为--conf keyvalue追加到spark-submit命令见 SparkStarter.java。生产模式YARN 集群 / 客户端模式把--master与--deploy-mode换成 YARN 即可在生产集群运行# Spark on YARN cluster mode ./bin/start-seatunnel-spark-3-connector-v2.sh --master yarn --deploy-mode cluster --config config/example.conf # Spark on YARN client mode ./bin/start-seatunnel-spark-3-connector-v2.sh --master yarn --deploy-mode client --config config/example.conf八、从示例走向真实作业从本文示例出发保留env块把FakeSource替换为真实 source 连接器选择 Source Connectors 并按对应文档配置参数把Console替换为目标 sink 连接器仅在源表结构与目标结构不一致时添加 transform运行前对照 Job Configuration Guide 的验证清单自查JAVA_HOME与 Java 版本、必需连接器插件是否安装、第三方驱动是否存在、source/sink 的凭据与网络连通性、目标表/主题/路径是否已存在、job.mode与所选连接器能力是否匹配若从源码树运行示例可参考seatunnel-examples/seatunnel-spark-connector-v2-example模块其入口类为org.apache.seatunnel.example.spark.v2.SeaTunnelApiExample详见 SeaTunnel With Spark想了解 SeaTunnel 自带默认引擎 Zeta可阅读 Quick Start With SeaTunnel Engine 以获得最短的本地验证路径。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考