ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

SeaTunnel Spark Translation Layer 深度解析:从 SeaTunnelSource 到 Spark Data Source V2 的适配原理与源码实践

SeaTunnel Spark Translation Layer 深度解析:从 SeaTunnelSource 到 Spark Data Source V2 的适配原理与源码实践 SeaTunnel Spark Translation Layer 深度解析从 SeaTunnelSource 到 Spark Data Source V2 的适配原理与源码实践【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 通过统一的连接器 APISeaTunnelSource/SeaTunnelSink/SeaTunnelTransform让同一套连接器可以运行在 Flink、Spark、SeaTunnel EngineZeta等多种执行引擎上而连接器作者只需要实现 SeaTunnel 自己的接口。本文聚焦其中的 Spark 专属路径讲解 SeaTunnel 如何把连接器语义重新解释为 Spark 原生概念——DataSource Reader、Input Partition、InternalRow、DataSource Writer 与 Commit Message并结合仓库源码剖析版本分治、Source/Sink 适配、Schema 与行转换、提交与恢复等核心机制。读完本文你将掌握 Spark 翻译层的整体架构、关键适配器类的调用链以及 schema 转换、commit/abort 等最容易出错的边界在哪里。为什么 Spark 翻译层不是换个接口名SeaTunnel 有独立的执行引擎Zeta也支持将同一套连接器跑在 Flink 与 Spark 之上。Spark 并不是与 Flink 同构的另一个引擎它的执行模型、数据源接口和提交生命周期都存在显著差异。因此 Spark 翻译层要做的不只是接口重命名而是把 SeaTunnel 的契约重新解释为 Spark 兼容的执行模型。具体差异体现在执行模型Spark 倾向于规划好的分区 Reader 执行模式而不是 Flink 那种持续活跃的 enumerator / runtime coordinator 模型数据源接口Spark 2.4 与 Spark 3.x 使用的 Data Source V2 API 并不相同提交生命周期Spark sink 有自己的一套 writer 与 commit message 模型需要与 SeaTunnel 的 writer / committer 语义桥接。设计目标与高层映射Spark 翻译层的核心目标是在把 SeaTunnel 语义适配到 Spark 原生概念datasource readers、input partitions、internal rows、datasource writers 和 commit messages的同时尽量让连接器作者免受 Spark 特定复杂度的打扰。关键不是让 Spark 看起来与 Flink 一致而是保持连接器语义的忠实性。从概念上看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 partition planning根据 SeaTunnel 的 split 信息规划 Spark 输入分区row and schema conversion在SeaTunnelRow/SeaTunnelDataType与 SparkInternalRow/StructType之间转换sink commit and abort behavior把 SeaTunnel 的提交语义映射到 Spark 的 writer / commit message 生命周期。版本分治Spark 2.4 与 Spark 3.x 的模块划分Spark 翻译在 SeaTunnel 中是版本感知的。Spark 2.4 与 Spark 3.x 暴露的 datasource API 并不相同因此 SeaTunnel 为各大版本分别维护独立的翻译模块与适配器。这一点意味着一个连接器在 SeaTunnel API 层行为正确并不代表底层不需要 Spark 版本特定的适配行为。仓库中的实际模块划分位于seatunnel-translation/seatunnel-translation-spark/seatunnel-translation-spark/ ├── seatunnel-translation-spark-common/ # 共享的类型转换、行转换、多表管理、工具类 ├── seatunnel-translation-spark-2.4/ # Spark 2.4 Data Source V2旧版 API └── seatunnel-translation-spark-3.3/ # Spark 3.x Data Source V2新版 TableProvider API从源码结构看Spark 2.4 模块seatunnel-translation-spark-2.4采用旧的DataSourceV2 / ReadSupport / MicroBatchReadSupport体系入口为SeaTunnelSourceSupport实现了DataSourceV2, ReadSupport, MicroBatchReadSupport, DataSourceRegistersource 侧按批/微批拆分为BatchSourceReader/MicroBatchSourceReader与CoordinatedBatchPartitionReader/ParallelBatchPartitionReader等sink 侧则有SparkSink、SparkDataSourceWriter、SparkDataWriterFactory、SparkStreamWriter、SparkWriterCommitMessage等类Spark 3.3 模块seatunnel-translation-spark-3.3采用新版TableProvider / Table体系入口为SeaTunnelSparkSourcesource与SeaTunnelSparkSinksink它们都实现DataSourceRegister, TableProvider并分别通过SeaTunnelSourceTable/SeaTunnelSinkTable暴露给 Spark 执行引擎公共模块seatunnel-translation-spark-common承载两个版本共享的逻辑TypeConverterUtils类型双向转换、InternalRowConverter/SeaTunnelRowConverter行双向转换、InternalRowCollector/InternalMultiRowCollector收集器适配、MultiTableManager多表元数据管理以及OffsetDateTimeUtils/InstantConverterUtils等时间处理工具。Source 侧适配从 split 到 Input Partition适配职责在 source 侧Spark 翻译层通常需要完成四件事以 Spark 期望的形式暴露 schema根据 SeaTunnel 的 split 信息规划分区为每个分区创建 reader把 SeaTunnel 输出转换为 SparkInternalRow。与 Flink 相比Spark 更强调规划好的分区与 reader 执行而非持续活跃的 enumerator / runtime coordinator 模型这一差异直接塑造了适配器设计。入口SeaTunnelSparkSourceSpark 3.xSpark 3.x 路径的入口是SeaTunnelSparkSource源码SeaTunnelSparkSource.java它实现了DataSourceRegister, TableProvidershortName()返回SeaTunnelSource用于 Spark SPI 发现inferSchema(...)返回nullsupportsExternalMetadata()返回true——SeaTunnel 不依赖 Spark 推断 schemaschema 由 SeaTunnel 自身提供getTable(...)返回SeaTunnelSourceTable再由SeaTunnelScan实现Scan与SeaTunnelBatch/SeaTunnelMicroBatch提供批式/微批式物理计划。分区规划planInputPartitions分区规划的核心逻辑在SeaTunnelBatch.planInputPartitions()源码SeaTunnelBatch.javaOverride public InputPartition[] planInputPartitions() { InputPartition[] partitions; if (source instanceof SupportCoordinate) { // 需要协调的 source只有一个 partition由协调者统一分配 split partitions new SeaTunnelBatchInputPartition[1]; partitions[0] new SeaTunnelBatchInputPartition(0); } else { // 普通 source按并行度创建 partition partitions new SeaTunnelBatchInputPartition[parallelism]; for (int partitionId 0; partitionId parallelism; partitionId) { partitions[partitionId] new SeaTunnelBatchInputPartition(partitionId); } } return partitions; }从源码可以看出两种读取模式Coordinated 模式当 source 实现了SupportCoordinate例如 Kafka、JDBC 等需要协调 split 分配的连接器时只规划 1 个 partition由内部的CoordinatedBatchPartitionReader统一调度Parallel 模式普通 source 按配置的并行度创建多个 partition每个 partition 对应一个ParallelBatchPartitionReader并行读取。Spark 2.4 侧的SeaTunnelSourceSupport.createReader(...)则会从EnvCommonOptions.PARALLELISM读取并行度默认 1并把getProducedCatalogTables()的结果包装成MultiTableManager后交给BatchSourceReader或MicroBatchSourceReader源码SeaTunnelSourceSupport.java。适配器结构示意以translation-layer.md中的示意代码为骨架Spark 侧适配器在概念上可以概括为三层。以下代码用于说明适配结构实际类名与生命周期以各版本源码为准// 示意SeaTunnelSource - Spark DataSourceReader public class SparkSourceT, SplitT extends SourceSplit, StateT implements DataSourceReader { private final SeaTunnelSourceT, SplitT, StateT seaTunnelSource; Override public StructType readSchema() { // 把 SeaTunnel schema 转换为 Spark StructType CatalogTable catalogTable seaTunnelSource.getProducedCatalogTables().get(0); return SparkTypeConverter.convert(catalogTable.getTableSchema()); } Override public ListInputPartitionInternalRow planInputPartitions() { // 创建 enumerator 并生成 splits每个 split 包装为一个 InputPartition SourceSplitEnumeratorSplitT, StateT enumerator seaTunnelSource.createEnumerator(new SparkEnumeratorContext()); enumerator.open(); enumerator.run(); ListSplitT splits collectAllSplits(enumerator); return splits.stream() .map(split - new SparkInputPartition(seaTunnelSource, split)) .collect(Collectors.toList()); } }// 示意每个 Spark InputPartition 负责为某个 split 创建分区 reader public class SparkInputPartitionT, SplitT extends SourceSplit implements InputPartitionInternalRow { private final SeaTunnelSourceT, SplitT, ? seaTunnelSource; private final SplitT split; Override public InputPartitionReaderInternalRow createPartitionReader() { SourceReaderT, SplitT seaTunnelReader seaTunnelSource.createReader(new SparkReaderContext()); return new SparkPartitionReader(seaTunnelReader, split); } }// 示意Spark 分区 reader 内部驱动 SeaTunnel reader并做逐条行转换 public class SparkPartitionReaderT, SplitT extends SourceSplit implements InputPartitionReaderInternalRow { private final SourceReaderT, SplitT seaTunnelReader; private final QueueInternalRow buffer new LinkedList(); Override public boolean next() throws IOException { if (!buffer.isEmpty()) { return true; } seaTunnelReader.pollNext(new CollectorT() { Override public void collect(T record) { InternalRow row SparkTypeConverter.convert(record); // 行转换 buffer.offer(row); } }); return !buffer.isEmpty(); } Override public InternalRow get() { return buffer.poll(); } }关于SeaTunnelSource接口本身的职责getBoundedness、createReader、createEnumerator、restoreEnumerator、getProducedCatalogTables、split/enumerator state 序列化器等可以参考 Source Architecture 中的完整接口定义与交互时序图。Sink 侧适配把 SeaTunnel 提交语义映射到 DataSourceWriter适配职责在 sink 侧Spark 翻译层把 SeaTunnel 的 sink 行为映射到 Spark 的 datasource writer 契约。典型职责包括创建 writer factoryDataWriterFactory把 executor 上的 commit message 带回 driver协调 commit 与 abort 路径把 SeaTunnel 的重试语义映射为 Spark 兼容行为。当 sink 不是简单的 append-only、需要幂等或事务性行为时这部分尤其关键。入口SeaTunnelSparkSinkSpark 3.xSpark 3.x 路径的入口是SeaTunnelSparkSink源码SeaTunnelSparkSink.javashortName()返回SeaTunnelSinkgetTable(...)返回SeaTunnelSinkTable。sink 侧从SeaTunnelWrite/SeaTunnelBatchWrite到SeaTunnelSparkDataWriterFactory/SeaTunnelSparkDataWriter/SeaTunnelSparkWriterCommitMessage构成了完整的写路径。Spark 2.4 侧则对应SparkSink、SparkDataSourceWriter、SparkDataWriterFactory、SparkStreamWriter、SparkWriterCommitMessage等类。适配器结构示意// 示意SeaTunnelSink - Spark DataSourceWriter public class SparkSinkIN, WriterStateT, CommitInfoT implements DataSourceWriter { private final SeaTunnelSinkIN, WriterStateT, CommitInfoT, ? seaTunnelSink; Override public DataWriterFactoryInternalRow createWriterFactory() { return new SparkDataWriterFactory(seaTunnelSink); } Override public boolean useCommitCoordinator() { // 只有当 sink 提供 committer 时才启用 commit coordinator return seaTunnelSink.createCommitter().isPresent(); } Override public void commit(WriterCommitMessage[] messages) { OptionalSinkCommitterCommitInfoT committerOpt seaTunnelSink.createCommitter(); if (committerOpt.isPresent()) { ListCommitInfoT commitInfos Arrays.stream(messages) .map(msg - ((SparkCommitMessageCommitInfoT) msg).getCommitInfo()) .collect(Collectors.toList()); ListCommitInfoT failed committerOpt.get().commit(commitInfos); if (!failed.isEmpty()) { throw new IOException(Some commits failed: failed); } } } Override public void abort(WriterCommitMessage[] messages) { OptionalSinkCommitterCommitInfoT committerOpt seaTunnelSink.createCommitter(); if (committerOpt.isPresent()) { ListCommitInfoT commitInfos Arrays.stream(messages) .map(msg - ((SparkCommitMessageCommitInfoT) msg).getCommitInfo()) .collect(Collectors.toList()); committerOpt.get().abort(commitInfos); } } }需要说明的是SeaTunnel 的SeaTunnelSink工厂接口本身提供了丰富的可选提交能力createWriter必选、createCommitter可选逐 writer 提交、createAggregatedCommitter可选全局聚合提交、以及对应的 state / commitInfo / aggregatedCommitInfo 序列化器。Spark 翻译层正是围绕这些钩子把两阶段提交映射到 Spark 的 commit message 生命周期。完整的接口定义与 XA 事务、文件原子重命名、Hive/Iceberg 表级提交等实现示例见 Sink Architecture。Schema 与行转换最敏感的边界Spark 翻译层高度依赖 schema 转换因为 Spark 通过自己强类型的 row 与 schema 模型执行。翻译层需要把CatalogTable/TableSchemaSeaTunnelDataTypeSeaTunnelRow映射为 Spark 概念StructTypeSpark SQL 数据类型InternalRow这是 Spark 路径中最敏感的边界之一尤其是decimals、timestamps、嵌套类型、nullability。类型映射TypeConverterUtils仓库中承担 SeaTunnel 与 Spark 类型双向转换的是TypeConverterUtils源码TypeConverterUtils.java。SeaTunnel → Spark 的核心映射如下SeaTunnelDataType (SqlType)Spark DataType备注NULLNullTypeSTRINGStringTypeBOOLEANBooleanTypeTINYINTByteTypeSMALLINTShortTypeINTIntegerTypeBIGINTLongTypeFLOATFloatTypeDOUBLEDoubleTypeBYTESBinaryTypeDATEDateTypeTIMELongType通过 Metadata 标记LOGICAL_TIME_TYPE_FLAG反转换时还原为LocalTimeTypeTIMESTAMPTimestampTypeTIMESTAMP_TZOffsetDateTimeUtils.OFFSET_DATETIME_WITH_DECIMAL带时区语义通过LOGICAL_TIMESTAMP_WITH_OFFSET_TYPE_FLAG标记DECIMALDecimalType(precision, scale)精度与标度完整保留ARRAYArrayType递归转换元素类型MAPMapType递归转换 key/value 类型ROWStructType递归转换每个字段支持嵌套反向转换Spark → SeaTunnel通过静态映射表TO_SEA_TUNNEL_TYPES加上对ArrayType/MapType/DecimalType/StructType的分支处理完成StructType转换时会读取字段 Metadata 中的LOGICAL_TIME_TYPE_FLAG与LOGICAL_TIMESTAMP_WITH_OFFSET_TYPE_FLAG从而把 Spark 的LongTypeTIME与特殊 decimal 类型还原为 SeaTunnel 的LOCAL_TIME_TYPE与OFFSET_DATE_TIME_TYPE。此外TypeConverterUtils.parcel(...)还会在多表/CDC 场景下为 schema 前置附加两个字段seatunnel_row_kindByteType表示 RowKind与seatunnel_table_idStringType表示表标识这是 SeaTunnel 在 Spark 上支撑多表与 CDC 数据的关键设计。行转换InternalRowConverter行级转换由InternalRowConverter源码InternalRowConverter.java完成它继承自 SeaTunnel 的RowConverterInternalRow负责在SeaTunnelRow与 SparkInternalRow之间双向转换内部使用 Spark 的SpecificInternalRow、MutableValue系列MutableInt、MutableLong、MutableByte等、ArrayBasedMapData、UTF8String、Decimal等 Catalyst 运行时对象。配套的SeaTunnelRowConverter负责反方向转换InternalRowCollector/InternalMultiRowCollector则把转换后的行送入下游收集器。在类型系统层面CatalogTable/TableSchema/Column/SeaTunnelDataType共同构成了 SeaTunnel 的便携式元数据与类型模型详见 Table Schema and Type System。该模型独立于具体引擎这是连接器能够在 Flink、Spark、Zeta 之间复用、并且 transform 与 sink 不需要各自实现引擎特定类型逻辑的前提。Commit 与恢复幂等、abort 与一致性Spark sink 执行有自己的一套 writer 与 commit message 模型。翻译层必须在桥接 SeaTunnel writer / committer 语义时不丢失以下保证幂等性期望idempotency expectationscommit 可能被调用多次SinkCommitter.commit()必须幂等失败处理行为failure handling behavior部分 writer 提交失败时应返回失败的 commitInfo 以便重试abort 正确性abort correctnesscheckpoint 失败时能正确回滚已 prepare 的事务连接器承诺的一致性保证consistency guarantees由连接器声明的 exactly-once / at-least-once 语义需要被完整传递。在 Spark 路径中useCommitCoordinator()的返回值决定了是否启用 Spark 的 commit coordinator当 sink 提供SinkCommitter时启用driver 侧在 checkpoint 成功后统一执行commit失败时走abort。SeaTunnel 还支持SinkAggregatedCommitter全局聚合提交适用于 Hive、Iceberg 等需要表级单次提交的场景。如果这个桥接薄弱用户通常会遇到三类问题重复副作用duplicate side effectscommit 重放导致数据重复abort 路径被破坏broken abort pathsprepare 之后无法回滚writer commit 不匹配writer commit mismatchescommit message 与 writer 状态对不上。关于 checkpoint 边界 两阶段提交如何实现 exactly-once以及与重试、幂等的关系可进一步阅读 Exactly-Once 和 Sink Architecture 中的失败与重试时序。常见故障点与排查思路Spark 翻译问题往往集中在以下区域schema 转换不匹配例如TIME在 Spark 中被表示为LongType、TIMESTAMP_TZ被表示为带元数据标记的特殊 decimal 类型如果 Metadata 标记丢失反向转换就会得到错误类型InternalRow 转换SeaTunnelRow与InternalRow的字段顺序、嵌套结构array/map/row处理、RowKind附加字段seatunnel_row_kind等都是容易出错的点datasource writer commit 行为commit 非幂等、prepare 阶段产生副作用、abort 未清理临时资源等Spark 2.4 与 Spark 3.x 适配器差异两个版本使用不同的 Data Source V2 API旧DataSourceV2/ReadSupport体系 vs 新版TableProvider/Table体系迁移连接器时必须分别验证。这些问题之所以隐蔽是因为故障表象可能看起来像连接器问题而真正的 bug 在翻译层内部。排查时建议先确认 schema 在转换前后是否完全一致尤其是 decimal 精度/标度、timestamp 语义、可空性、嵌套类型再检查InternalRow的字段序与元数据标记最后核对 commit/abort 路径的幂等性。推荐阅读路径围绕 Spark 翻译层可以按以下顺序深入Translation Layer——跨引擎翻译层总览先建立 Flink / Spark / Zeta 的整体视图本文——Spark 专属路径Source Architecture——SeaTunnelSource/SourceSplitEnumerator/SourceReader的接口契约与交互时序Sink Architecture——SeaTunnelSink/SinkWriter/SinkCommitter/SinkAggregatedCommitter的提交语义Table Schema and Type System——CatalogTable/TableSchema/SeaTunnelDataType的类型系统全貌。若需要结合源码继续阅读重点入口为seatunnel-translation/seatunnel-translation-spark/目录下的三个模块核心类包括SeaTunnelSparkSource/SeaTunnelSparkSinkSpark 3.x 入口、SeaTunnelSourceSupport/SparkSinkSpark 2.4 入口、SeaTunnelBatch/SeaTunnelMicroBatch分区规划、SeaTunnelBatchPartitionReader/SeaTunnelInputPartitionReader分区读取以及公共模块中的TypeConverterUtils/InternalRowConverter/SeaTunnelRowConverter类型与行转换。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表