
SeaTunnel With Spark在现有 Spark 集群上运行 SeaTunnel 作业的完整指南【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunneloutput文章SeaTunnel With Spark在现有 Spark 集群上运行 SeaTunnel 作业的完整指南本文是 SeaTunnel 项目 Spark 引擎接入的官方指南。当你所在团队已经运营 Spark 集群、且主要负载是批量batch或混合型工作负载时通过spark.前缀配置、Spark 专属启动脚本与 Spark Translation Layer翻译层可以让 SeaTunnel 的 Source/Transform/Sink 作业无缝运行在 Spark 之上。读完本文你将掌握何时选择 Spark 引擎、如何编写 Spark 专属配置、如何用start-seatunnel-spark-3-connector-v2.sh在 YARN cluster/client 模式下提交作业、如何基于仓库源码运行官方示例以及 Spark 翻译层背后的工作原理。何时选择 Spark 引擎SeaTunnel 是一个多引擎支持的数据集成工具除了内置的 SeaTunnel EngineZeta外还支持 Apache Spark 与 Apache Flink 作为执行引擎。选择 Spark 的典型前提是你的组织已经在生产环境运行 Spark 集群希望复用现有运维体系周边工作负载以批处理batch为主你希望 SeaTunnel 作业与现有 Spark 生态Hive、HDFS、YARN、Kubernetes 等及部署模型对齐。需要注意的是官方在 Engine Overview 中的推荐路径是新项目、无现有 Flink/Spark 基础设施时优先使用 SeaTunnel EngineZeta只有当你已运营 Spark 集群并希望复用时才选择 Spark。此外从引擎能力对比表可以明确看到Spark 引擎不支持 CDC 连接器与 Schema Evolution因此在 CDC 实时同步、多表迁移等场景下应优先选择 SeaTunnel Engine 或 Flink。引擎能力对比摘录自 Engine Overview特性SeaTunnel EngineFlinkSpark批处理✅✅✅流处理✅✅✅CDC 支持✅✅❌Exactly-Once✅✅✅Schema Evolution✅✅❌REST APISeaTunnel 自身✅❌❌Standalone / Cluster 模式✅ / ✅✅ / ✅✅ / ✅Spark 专属配置env 块与spark.前缀SeaTunnel 的作业配置文件HOCON 格式中所有 Spark 专属选项都位于env块内并使用spark.前缀。这些配置会在提交阶段被解析并转换为spark-submit --conf参数。官方给出的 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 }常见参数含义与建议配置项含义备注spark.app.nameSpark 应用名称会在 YARN 与 Spark UI 中显示若不设置代码会回退到job.name再回退到默认值见 SparkRuntimeEnvironment.javaspark.sql.catalogImplementationSpark SQL 目录实现hive表示启用 Hive 元数据与 Hive 数据源搭配时使用spark.executor.memory每个 Executor 的内存如2gspark.executor.instancesExecutor 数量如2spark.yarn.priorityYARN 上的应用优先级如100spark.dynamicAllocation.enabled是否启用动态资源分配批量作业常显式关闭以固定资源底层解析机制从源码看env块中所有键值对都会被逐一取出并原样写入SparkConf最终由 SparkStarter.buildFinal() 转换为spark-submit --conf keyvalue参数。具体逻辑见 SparkStarter.getSparkConf()static MapString, String getSparkConf(String configFile, ListString variables) { Config appConfig ConfigBuilder.of(configFile, variables); return appConfig.getConfig(env).entrySet().stream() .collect(Collectors.toMap( Map.Entry::getKey, e - e.getValue().unwrapped().toString())); }注意env块中的全部键值对都会进入 SparkConf包括非spark.前缀的通用项如parallelism而SparkRuntimeEnvironment.createSparkConf()会把这些配置逐一sparkConf.set(key, value)因此不要把与 Spark 无关的敏感或临时配置放进env块以免被透传到 Spark 侧。此外SparkRuntimeEnvironment还有两个值得了解的行为见 SparkRuntimeEnvironment.java自动启用 Hive 支持当作业配置中的 source 或 sink 插件名包含hive不区分大小写时会自动调用SparkSession.builder().enableHiveSupport()见checkIsContainHive。流式批处理时长默认 Spark Streaming 批处理时长为 5 秒可通过spark.stream.batchDuration覆盖。命令行示例YARN 集群模式与客户端模式Spark 引擎通过专属启动脚本提交作业脚本名对应 Spark 大版本Spark 3.x./bin/start-seatunnel-spark-3-connector-v2.shSpark 2.4.x./bin/start-seatunnel-spark-2-connector-v2.shSpark on YARN cluster 模式./bin/start-seatunnel-spark-3-connector-v2.sh --master yarn --deploy-mode cluster --config config/example.confSpark on YARN client 模式./bin/start-seatunnel-spark-3-connector-v2.sh --master yarn --deploy-mode client --config config/example.conf命令行参数说明--master与--deploy-mode的定义见 SparkCommandArgs.java--master-mSpark master 地址支持spark://host:port、mesos://host:port、yarn、k8s://https://host:port、local默认local[*]--deploy-mode-e仅支持cluster与client两种默认client传入其他值会抛出IllegalArgumentException--configSeaTunnel 作业配置文件路径--check-c仅校验配置不实际执行--encrypt/--decrypt配置文件加解密。底层发生了什么启动脚本 → spark-submitSparkStarter的职责是把 SeaTunnel 作业生成一条完整的spark-submit命令见 SparkStarter.java其中关键步骤包括解析env块得到 SparkConf收集lib/下的依赖 jar、根据作业配置中实际使用的 source/transform/sink 插件从connectors/目录解析出对应插件 jar见getConnectorJarDependencies通过SeaTunnelSourcePluginDiscovery/SeaTunnelSinkPluginDiscovery按plugin_name精确匹配拼装--class org.apache.seatunnel.core.starter.spark.SeaTunnelSpark、--master、--deploy-mode、--jars、--conf、--name等参数cluster 模式下还会先把整个plugins/目录打包成 tar.gz 并通过--files分发到 Executor同时把配置文件一并分发见ClusterModeSparkStarter.buildCommands()这是 cluster 模式下插件能在远端 Executor 加载的关键。作业真正的入口类是 SeaTunnelSpark.main()它解析命令行后调用SeaTunnel.run(...)进入任务执行阶段。最小可运行作业示例下面的示例在 Spark 上运行FakeSource生成 16 行测试数据 →FieldMapper做字段重命名 →Console打印到控制台。env { parallelism 1 spark.app.name example spark.sql.catalogImplementation hive spark.executor.memory 2g spark.executor.instances 1 spark.yarn.priority 100 spark.dynamicAllocation.enabled false } 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 } }配置要点插件数据流接线plugin_output声明上游产出的数据集标识下游插件的plugin_input引用该标识形成FakeSource(fake) → FieldMapper(fake→fake1) → Console(fake1)的链式结构FakeSourcerow.num控制生成行数schema.fields声明字段名与类型也可通过schema.fields.name { type string, max 100 }等方式配置随机值范围详见 FakeSource 文档实际文件名为fake.md可通过docs/en/connectors/source/目录确认FieldMapperfield_mapper的 key 是原字段名value 是映射后的新字段名示例将name重命名为new_nameConsole Sink把数据打印到终端适合联调验证。运行该作业后控制台会输出类似下面的日志字段、类型与 16 行随机数据fields : name, age types : STRING, INT row1 : elWaB, 1984352560 row2 : uAtnp, 762961563 ... row16 : SGZCr, 94186144更多 transform 选项可参考 Transforms Catalog 与 Transform Common Options。执行链路Source → Transform → Sink在 Spark 引擎上作业按“先 Source、再 Transform、最后 Sink”的顺序串联执行数据以DatasetTableInfo形式在各阶段之间传递见 SparkExecution.execute()datasets sourcePluginExecuteProcessor.execute(datasets); datasets transformPluginExecuteProcessor.execute(datasets); sinkPluginExecuteProcessor.execute(datasets);三个 ExecuteProcessorSourceExecuteProcessor、TransformExecuteProcessor、SinkExecuteProcessor分别负责把 SeaTunnel 插件接入 Spark 的 Dataset/DataFrame 计算管线。在源码仓库中运行示例如果你的环境是从源码检出source checkout可直接使用仓库自带的示例模块示例模块seatunnel-examples/seatunnel-spark-connector-v2-example入口类org.apache.seatunnel.example.spark.v2.SeaTunnelApiExample入口类的实现见 SeaTunnelApiExample.java其默认读取/examples/spark.batch.conf作为配置也可以在运行时通过第一个参数传入自定义配置文件路径java -cp seatunnel-spark-connector-v2-example-classpath \ org.apache.seatunnel.example.spark.v2.SeaTunnelApiExample \ /path/to/your/spark.conf生产环境部署时建议按照 Deployment 文档 下载并解压发行包并确保Spark 版本 2.4.0当前仓库的 Spark starter 支持 2.4 与 3.x 两个大版本线见seatunnel-core/seatunnel-spark-starter/下seatunnel-spark-2-starter、seatunnel-spark-3-starter等模块在${SEATUNNEL_HOME}/config/seatunnel-env.sh中设置SPARK_HOME指向 Spark 部署目录见 seatunnel-env.shSPARK_HOME${SPARK_HOME:-/opt/spark}将作业配置写入config/目录如基于config/v2.batch.config.template或config/v2.streaming.conf.template修改。深入原理Spark Translation Layer翻译层SeaTunnel 通过Translation Layer把自身插件 API 适配到 Spark 的执行模型相关设计与实现集中在seatunnel-translation/seatunnel-translation-spark/seatunnel-translation-spark-common/seatunnel-translation-spark-2.4/seatunnel-translation-spark-3.3/为什么 Spark 翻译是特殊的Spark 的执行模型、数据源接口与提交生命周期和 Flink 有本质差异因此翻译层不能只是“接口改名”而必须把 SeaTunnel 的语义重新解释为 Spark 兼容的执行模型。其核心目标是在适配 Spark 原生概念datasource reader、input partition、InternalRow、datasource writer 与 commit message的同时尽量保持 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翻译的关键点集中在三处source 分区规划、行与 schema 转换、sink 提交与回滚行为。完整设计见 Spark Translation Layer 文档。Source 侧适配Spark 端更强调“预规划的分区 按分区执行 Reader”而非 Flink 那种持续活跃的 enumerator/runtime coordinator 模型。Spark 翻译层在 source 侧需要以 Spark 期望的形态暴露 schema根据 SeaTunnel 的 split 信息规划 input partition为每个分区创建 reader把 SeaTunnel 输出转换为 SparkInternalRow。对应实现可查看 Spark 3.3 模块下的 source/partition/batch 与 source/partition/micro 目录批量场景使用SeaTunnelBatch/SeaTunnelBatchPartitionReaderFactory微批micro-batch场景使用SeaTunnelMicroBatch/CoordinatedMicroBatchPartitionReader等类。Sink 侧适配翻译层把 SeaTunnel sink 行为映射到 Spark 的 datasource writer 契约典型职责包括创建 writer factory从 Executor 回传 commit message协调 commit 与 abort 路径把重试语义映射为 Spark 兼容行为。当 sink 不是纯 append 型、需要幂等或事务语义时这部分尤为关键。可参考 sink/write 目录下的SeaTunnelWrite、SeaTunnelSparkDataWriter、SeaTunnelSparkDataWriterFactory、SeaTunnelSparkWriterCommitMessage等类。Schema 与行转换Spark 依赖自身强类型的行与 schema 模型执行因此翻译层需要把 SeaTunnel 侧的CatalogTable/TableSchema、SeaTunnelDataType、SeaTunnelRow映射为 Spark 的StructType、Spark SQL 数据类型与InternalRow。这是 Spark 路径上最敏感的边界之一尤其是decimal、timestamp、嵌套类型与可空性nullability的转换。相关实现见 serialization 目录下的SeaTunnelRowConverter、InternalRowConverter与 utils/TypeConverterUtils.java。版本分裂2.4 与 3.xSpark 2.4 与 Spark 3.x 的数据源 API 并不完全相同因此 SeaTunnel 为两条大版本线分别维护独立的翻译模块与适配器seatunnel-translation-spark-2.4与seatunnel-translation-spark-3.3。这意味着一个连接器在 SeaTunnel API 层行为正确底层仍可能需要 Spark 版本专属的适配行为。提交与恢复的坑Spark sink 执行有自己的一套 writer / commit message 模型翻译层必须在桥接时保留幂等性预期、失败处理、abort 正确性以及连接器承诺的一致性语义。如果桥接薄弱用户通常会遇到重复副作用duplicate side effects、abort 路径损坏、writer commit 不匹配。常见问题集中在schema 转换不匹配、InternalRow转换错误、datasource writer 提交行为异常、Spark 2.4 与 3.x 适配差异。这类问题很隐蔽——表面看像连接器 bug实际根源在翻译层。与 SeaTunnel EngineZeta的取舍若你并非必须运行在 Spark 上官方建议优先评估内置的 SeaTunnel Engine它开箱即用、无 Zookeeper/HDFS 等外部依赖、资源开销更低、启动更快且对 CDC、多表同步等同步类负载支持更完整。下表总结了选择依据来自 Engine Overview场景推荐引擎无大数据基础设施的新项目SeaTunnel EngineCDC 与实时同步SeaTunnel Engine已有 Flink 基础设施Flink已有 Spark 基础设施、大规模批量 ETLSpark低资源环境SeaTunnel Engine复杂流处理Flink需要特别提醒的是在 Spark 引擎上运行 SeaTunnel 时作业的提交与监控应通过 Spark 自身的工具链完成如spark-submit、Spark History Server因为 SeaTunnel 的 REST API V2 只由 SeaTunnel EngineZeta服务端实现Flink/Spark 引擎下不可用。下一步学习路径Quick Start With Spark完整的分步上手部署、配置 SPARK_HOME、运行示例、观察输出Spark Translation Layer翻译层架构细节Job Configuration Guide作业配置整体规范Engine Overview三大引擎能力对比与选型SeaTunnel Engine与默认引擎做对比时参考。/output文章【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考