SeaTunnel 在 Spark 上的快速入门:部署、配置与运行你的第一个同步作业
2026/9/18 13:17:39 网站建设 项目流程

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 Engine(Zeta)服务集群,这是 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. 环境准备

  • 安装 Java(8 或 11,高于 8 的版本理论上也可用),并正确设置JAVA_HOME

2. 下载二进制发行包

从 SeaTunnel 下载页获取seatunnel-<version>-bin.tar.gz,或通过终端下载解压:

export version="3.0.0" wget "https://archive.apache.org/dist/seatunnel/${version}/apache-seatunnel-${version}-bin.tar.gz" tar -xzvf "apache-seatunnel-${version}-bin.tar.gz"

Windows 用户可下载.zip包后用文件管理器或 PowerShell 解压。

3. 安装连接器插件

自 2.2.0-beta 起,二进制包默认不再内置连接器依赖,首次使用需执行插件安装脚本:

sh bin/install-plugin.sh

关键细节:

  • 不需要安装全部连接器。通过修改config/plugin_config即可指定所需插件。以本文示例作业为例,只需connector-fakeconnector-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 与校验和(需curlmktemp以及sha512sum/sha1sum/shasum/openssl之一);也可通过环境变量SEATUNNEL_MAVEN_REPOSITORY指向 HTTPS Maven 镜像,或用SEATUNNEL_PLUGIN_DOWNLOAD_METHOD=maven保留 Mavensettings.xml的镜像、认证仓库、代理等行为。

三、Step 2:部署并配置 Spark

1. 下载 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_HOME(Flink 模式使用)以及 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 行模拟数据,字段name(STRING)与age(INT)
transform对在途数据做变换(可选)FieldMapper重命名字段,将name映射为new_name
sink定义数据去向Console将结果打印到控制台

关键参数说明

  • plugin_output/plugin_input:这是理解 SeaTunnel 数据流的核心约定。plugin_output为 source/transform 产出的数据流命名,plugin_input告诉下游 transform/sink 消费哪条上游流。本示例中FakeSource输出名为fakeFieldMapper消费fake并输出fake1Console消费fake1。当作业只有一个上游路径时,SeaTunnel 通常能按默认约定自动串联,但显式命名更利于多源、多分支作业的可读性,详见 Job Configuration Guide。
  • row.num:FakeSource 生成的行数,定义于 FakeSourceOptions.java,本示例为 16。
  • schema.fields:声明字段名与类型的映射,namestringageint,对应输出日志中的types : STRING, INT
  • field_mapper:FieldMapper 的字段重命名映射,age = age保持不变,name = new_name将原字段name重命名为new_name。因此最终控制台输出仍只有两列(重命名后的new_nameage)。
  • job.modeBATCHSTREAMING。注意仓库自带的v2.streaming.conf.template是流式示例(job.mode = "STREAMING"checkpoint.interval = 2000),与本文批处理示例在env块上不同,运行前请按需调整。

更完整的配置概念请查阅 Config Concept;变换类插件的通用参数见 Transform Common Options。

五、Step 4:运行 SeaTunnel 应用

根据 Spark 主版本选择对应启动脚本:

Spark 2.4.x

cd "apache-seatunnel-${version}" ./bin/start-seatunnel-spark-2-connector-v2.sh \ --master local[4] \ --deploy-mode client \ --config ./config/v2.streaming.conf.template

Spark 3.x.x

cd "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 核运行;生产环境可替换为yarnspark://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_MAIN="org.apache.seatunnel.core.starter.spark.SparkStarter"APP_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 row=1 : elWaB, 1984352560 row=2 : uAtnp, 762961563 row=3 : TQEIB, 2042675010 row=4 : DcFjo, 593971283 row=5 : SenEb, 2099913608 row=6 : DHjkg, 1928005856 row=7 : eScCM, 526029657 row=8 : sgOeE, 600878991 row=9 : gwdvw, 1951126920 row=10 : nSiKE, 488708928 row=11 : xubpl, 1420202810 row=12 : rHZqb, 331185742 row=13 : rciGD, 1112878259 row=14 : qLhdI, 1457046294 row=15 : ZTkRx, 1240668386 row=16 : SGZCr, 94186144

日志解读:

  • fields : name, agetypes : STRING, INT来自作业schema定义,与schema.fields完全对应;
  • 16行记录(row=1row=16),与row.num = 16一致;
  • 每行由两个字段构成,经FieldMapper重命名后输出列语义不变(字段名变为new_nameage),随机内容来自 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/:公共部分,如InternalRowConverterSeaTunnelRowConverterTypeConverterUtils等行列与类型转换工具;
  • seatunnel-translation/seatunnel-translation-spark/seatunnel-translation-spark-2.4/:Spark 2.4 适配(如SparkSinkSparkDataSourceWriterSeaTunnelInputPartitionReader等);
  • seatunnel-translation/seatunnel-translation-spark/seatunnel-translation-spark-3.3/:Spark 3.3 适配(如SeaTunnelSparkSourceSeaTunnelBatch/SeaTunnelMicroBatchSeaTunnelSparkSinkSeaTunnelWrite等)。

翻译层最敏感的边界包括:source 侧的分区规划(从 SeaTunnel split 信息映射为 Spark 输入分区)、schema 与行的转换(CatalogTable/TableSchemaSeaTunnelDataTypeSeaTunnelRowStructType、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 key=value追加到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),仅供参考

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

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

立即咨询