十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

深入解析 Delta Spark V2 Connector:基于 Delta Kernel 的 Spark 数据源 V2 集成架构

深入解析 Delta Spark V2 Connector:基于 Delta Kernel 的 Spark 数据源 V2 集成架构 深入解析 Delta Spark V2 Connector基于 Delta Kernel 的 Spark 数据源 V2 集成架构【免费下载链接】deltaAn open-source storage framework that enables building a Lakehouse architecture with compute engines including Spark, PrestoDB, Flink, Trino, and Hive and APIs项目地址: https://gitcode.com/GitHub_Trending/del/deltaDelta Spark V2 Connector 是当前仓库中用于将 Delta Kernel 接入 Apache Spark 查询执行管线的新一代连接器。它基于 Spark 的 DataSource V2DSV2API 实现让 Spark 可以借助 Delta Kernel 读取 Delta 表并在此之上支持过滤器下推、数据跳过file skipping、列裁剪、元数据列、CDC 与流式读取等能力。阅读本文后你将掌握该连接器在 Spark 与 Delta Kernel 之间的桥梁定位、端到端的读取流程以及其核心模块的源码级实现细节可直接对照仓库代码进行二次开发或问题排查。一、连接器定位为什么需要 Spark 与 Delta Kernel 之间的桥Delta Kernel 是 Delta Lake 的独立读端内核它不依赖 Spark 运行环境负责从 Catalog 和 Delta 事务日志_delta_log中解析表的当前状态。而 Spark 生态中的查询执行、任务调度、Parquet 扫描等能力则完全由 Spark 引擎本身提供。两者之间需要一个适配层把 Spark 的读表请求翻译成 Delta Kernel 能够理解的语言再把 Kernel 返回的文件列表交还给 Spark 去执行实际扫描。Delta Spark V2 Connector代码位于 spark/v2 目录正是这个适配层。spark/v2/README.md 明确给出了它的核心设计通过 Spark 的DataSource V2 API与 Spark 查询执行管线集成让 Spark 借助Delta Kernel读取 Delta 表由连接器负责请求表 Schema、下推过滤器、获取待扫描文件列表。从源码结构看连接器的核心实现集中在io.delta.spark.internal.v2包下包含catalog表实现、read扫描与读取、write写入、snapshot快照管理、kernelKernel 引擎等子模块并在java-shims与scala-shims中针对不同 Spark 版本4.0 / 4.1 / 4.2提供适配实现。二、高层设计一次查询的完整数据流spark/v2/README.md将连接器描述为位于 Spark 与 Delta Kernel 之间的桥梁并给出了四个关键环节Spark Driver 请求表 Schema并通过连接器下推静态static与动态dynamic过滤器连接器将请求翻译为 Delta Kernel API 调用向 Kernel 请求表 Schema、下推过滤器、获取待扫描文件列表Delta Kernel 解析表状态从 Catalog 与 Delta 日志中解析快照应用文件跳过逻辑file skipping返回必要的文件Spark Engine 进行分区与扫描以默认 128MB 的拆分粒度对文件分区使用 Spark 既有的ParquetFileFormat执行实际的 Parquet 扫描。这张架构图直观展示了数据流方向绿色箭头代表 Spark Driver 与连接器之间的交互请求 Schema、下推静态与动态过滤器紫色箭头代表连接器与 Delta Kernel 之间的交互请求 Schema、下推过滤器并获取扫描文件蓝色箭头则代表扫描文件最终回到 Spark由 Spark 的 Parquet Reader 完成实际数据读取。值得强调的是Delta Kernel 在这里承担了决策者的角色——它决定读取哪些文件而 Spark 承担了执行者的角色——它负责把文件拆分成任务并真正扫描数据。连接器则保证了两种角色之间的协议转换。三、表实现DeltaV2Table 与 DSV2 能力矩阵连接器的表入口是 DeltaV2Table.java它实现了 Spark DSV2 的多个关键接口DeltaV2Table.java#L92-L99Table表的基本描述Schema、分区、属性SupportsRead批量与微批读取SupportsWrite批量与流式写入SupportsMetadataColumns暴露_metadata元数据列SupportsRowLevelOperations行级操作如 DELETE 的规划入口。其能力集合定义在buildCapabilities()中DeltaV2Table.java#L108-L122EnumSet.of( TableCapability.BATCH_READ, // 批量读取 TableCapability.MICRO_BATCH_READ, // 微批读取流式 TableCapability.BATCH_WRITE, // 批量写入 TableCapability.STREAMING_WRITE); // 流式写入此外还会通过 Spark 版本 shim 附加 Schema 演进schema evolution能力。一个值得注意的细节是当表被固定到某个 time travel 版本时capabilities()会退化为只读集合BATCH_READMICRO_BATCH_READ因为时间旅行固定的表不可写入此时写操作会回退到 V1 路径DeltaV2Table.java#L454-L459。3.1 构造路径与选项合并DeltaV2Table支持两类构造方式对应两种使用场景基于路径直接传入文件系统路径如delta.\/path/to/table无 Catalog 表元数据基于 Catalog 表传入 Spark 的CatalogTable从中提取表位置location与存储属性。构造时会把Catalog 表存储选项与用户选项合并用户选项优先级更高随后用合并后的选项构建 Hadoop 配置newHadoopConfWithOptions、创建 Kernel 引擎、并通过SnapshotManagerFactory创建快照管理器加载初始快照DeltaV2Table.java#L226-L289。3.2 Schema 的懒计算Schema 相关元数据公开 Schema、数据列 Schema、分区列 Schema、分区 Transform由内部类SchemaProvider在首次访问时懒计算并缓存DeltaV2Table.java#L643-L732。实现要点包括公开 Schema 会剥离 Delta 内部元数据removeInternalDeltaMetadata与removeInternalWriterMetadata分区列严格按partColNames指定的顺序排列该顺序可能与快照 Schema 中的顺序不一致必须保持以维持正确的分区行为分区列在公开 Spark Schema 中追加在数据列之后符合 Spark 的约定。四、Kernel 引擎与快照管理4.1 Kernel 引擎创建连接器通过 KernelEngineFactory.scala 创建默认的 Kernel 引擎KernelEngineFactory.scala#L24-L29def createDefaultEngine(hadoopConf: Configuration): KernelEngine { KernelDefaultEngine.create(hadoopConf) }即直接使用io.delta.kernel.defaults.engine.DefaultEngine以 Hadoop 配置为参数构建。而 KernelContext.scala 则负责 Kernel 资源的表作用域管理KernelContext.scala#L33-L50保留会话无关的文件系统选项session-invariant filesystem options与 LogStore在真正发生文件系统 I/O 时才把选项绑定到当前活跃的 Spark 会话避免可复用的连接器状态持有会话派生的设置与凭据其 LogStore 与懒创建的 Engine 跨 Spark 会话保持稳定。这种延迟物化 Hadoop 配置的设计是为了让连接器状态能够在多个 Spark 会话之间安全复用同时不泄漏任一会话的凭据。4.2 快照管理器路径表与 Unity Catalog 托管表快照加载由 SnapshotManagerFactory.java 根据表类型分派SnapshotManagerFactory.java#L58-L60路径型表→PathBasedSnapshotManagerUnity Catalog 托管表→UCManagedTableSnapshotManager使用UCCatalogManagedClient与UCClient处理协调提交。两者都实现DeltaV2SnapshotManager接口DeltaV2SnapshotManager.scala#L39-L95提供统一的操作面方法作用loadLatestSnapshot()加载最新快照loadSnapshotAt(version)加载指定版本快照版本号 ≥ 0getActiveCommitAtTime(timestampMillis, ...)按时间戳UTC 毫秒解析当时活跃的提交支持返回最后一次提交仅考虑可重建提交返回最早提交等策略checkVersionExists(version, ...)校验版本是否存在且可访问不可用时抛出VersionNotFoundExceptiongetTableChanges(engine, startVersion, endVersion)获取起始版本到结束版本之间的表变更CommitRange供流式读取使用在DeltaV2Table构造时time travel 固定版本通过checkVersionExists校验后由loadSnapshotAtCheckedVersion加载时间戳形式的 time travel 则先调用getActiveCommitAtTime解析出版本号DeltaV2Table.java#L378-L414。若请求时间早于最早可用提交或晚于最新提交会抛出TimestampOutOfRangeException。Kernel 抛出的TableNotFoundException也会被包装为连接器自己的异常类型避免 Catalog / 互操作层暴露 Kernel 内部类型。五、扫描构建过滤器下推、列裁剪与 LIMIT 下推扫描规划的核心是 DeltaV2ScanBuilder.scala它同时实现了多个 Spark DSV2 下推接口SupportsPushDownRequiredColumns列裁剪SupportsPushDownCatalystFiltersCatalyst 表达式过滤器下推SupportsPushDownLimitLIMIT 下推。5.1 过滤器下推pushFilters会把 Spark 传入的 Catalyst 表达式按分区列拆分为分区过滤器与数据过滤器DeltaV2ScanBuilder.scala#L106-L117分区过滤器是精确的下推后无需再次求值数据过滤器依赖 min/max 统计做文件级跳过data skipping并非行级精确因此会保留扫描后残余过滤器post-scan residual需要在扫描后由 Spark 重新求值。build()阶段会把这些 Catalyst 过滤器交给DeltaV2Snapshot.filesForScan走 V1 的数据跳过路径选出文件DeltaV2ScanBuilder.scala#L156-L220。这正是 README 中连接器下推过滤器、Kernel 应用文件跳过并返回文件列表这一环节的落地实现——文件选择复用 V1 的filesForScan而 Kernel 负责快照解析。5.2 LIMIT 下推pushLimit接受 Spark 优化器下推的 LIMIT 提示DeltaV2ScanBuilder.scala#L145-L151。实现把 LIMIT 视为尽力而为的提示当不存在残余过滤器时文件选择走 V1 的 limit-awarefilesForScan按记录数累计决定停止添加文件。由于剪枝发生在文件粒度单个文件可能包含超过 LIMIT 的行数例如LIMIT 5但单个文件有 1000 行所以isPartiallyPushed()保持默认的trueSpark 会保留 LIMIT 算子作为兜底确保最终结果行数精确。5.3 列裁剪pruneColumns会剔除分区列与 CDC 注入列仅保留真正需要的数据列DeltaV2ScanBuilder.scala#L122-L131CDC 列在后续由CDCReadFunction注入。此外DeltaV2ScanBuilder还提供了forUnfilteredScan适配器供当前尚未实现过滤器下推的调用方如 Delta Sharing 的 DSv2做全表无过滤扫描。5.4 任务拆分与列式执行文件选定后Spark Engine 按默认 128MB 拆分为FilePartition并调度到执行端。DeltaV2ReaderFactory.java 实现了PartitionReaderFactory根据supportsColumnar标志决定返回行式还是列式ColumnarBatch读取器两者最终都落到DeltaV2PartitionReader由 Spark 既有的ParquetFileFormat执行 Parquet 扫描。列式读取路径与DeltaParquetFileFormatV2、ColumnVectorWithFilter等组件配合服务于列式批量读取场景。六、元数据列与行级跟踪Row TrackingDeltaV2Table通过SupportsMetadataColumns暴露一个名为_metadata的元数据列DeltaV2Table.java#L486-L523其结构与 Spark 文件源的BASE_METADATA_FIELDS保持一致基础字段file_path、file_name、file_size、file_block_start、file_block_length、file_modification_time当表启用了行跟踪Row Tracking时额外追加row_id与row_commit_version两个 Long 字段。字段取值遵循PartitionedFile语义字段顺序与 V1 Delta 保持一致以维持对等但解析按名称进行顺序不参与正确性判定。行级操作SupportsRowLevelOperations的规划则通过DeltaRowLevelOperationBuilder接入 Delta 的 copy-on-write 操作框架。七、写入、流式与 CDC 支持连接器不止读还通过DeltaV2WriteBuilder、DeltaV2BatchWrite、DeltaV2StreamingWrite、DeltaV2DataWriter等组件提供批量写入与流式写入能力见 write 目录。读取侧还包含流式读取DeltaV2MicroBatchStream基于快照管理器提供的 commit-range 变更实现微批流CDC 读取CDCReadFunction、CDCSchemaContext与CDCDataFile支持 Change Data Feed删除向量DeletionVectorReadFunction配合ColumnVectorWithFilter处理带删除向量的文件。针对 Spark 4.2 还提供了DeltaV2Changelog系列DeltaV2ChangelogScanBuilder、DeltaV2ChangelogBatch等用于读取变更日志以及元数据-only DELETE 执行器DeltaMetadataOnlyDeleteExecutor。这些能力共同支撑了从纯读取到完整读写、从批量到流式的覆盖范围。八、测试体系与验证路径连接器配备了与实现一一对应的测试套件可作为理解行为与验证机制的入口读取侧V2ReadTest.java、DeltaV2ScanBuilderTest.java、V2LimitPushdownTest.java、V2MetadataReadTest.java流式读取V2StreamingReadTest.java、V2StreamingReadDistributedTest.java、DeltaV2MicroBatchStreamTest.java写入侧V2WriteTest.java、DeltaV2WriteTest.java、V2DDLTest.java快照与 Kernel 上下文PathBasedSnapshotManagerTest.java、KernelContextSuite.scala行跟踪V2RowTrackingReadTest.java列式/向量化读取DeltaV2ScanTest.java。如需在仓库内运行这些测试可参照根目录的 run-tests.py 与 spark/v2 模块的测试组织方式按对应 Spark 版本执行。九、使用边界与注意事项结合仓库现状使用该连接器时有几点值得注意依赖 Delta Kernel 与 Spark 4.x 的 DSV2 接口连接器按 Spark 4.0 / 4.1 / 4.2 分别提供 shimjava-shims、scala-shims不同 Spark 版本的行为差异由 shim 层吸收数据过滤器不是行级精确的min/max 统计跳过只在文件粒度生效残余过滤由 Spark 端保证正确性time travel 只读固定到历史版本的表只允许读取写操作会回退到 V1 路径LIMIT 下推是文件粒度的近似最终行数精确性依赖 Spark 保留的 LIMIT 算子兜底Delta Sharing 等部分调用方暂未接入过滤器下推forUnfilteredScan的存在说明这类场景目前走全表扫描。总体而言Delta Spark V2 Connector 通过 DSV2 接口把 Delta Kernel 的元数据解析与文件选择能力无缝嵌入 Spark 的执行管线既复用了 Spark 成熟的 Parquet 扫描与调度能力又将 Delta 表状态管理交给独立的 Kernel 内核是当前仓库中连接Spark 生态与Delta Kernel的关键桥梁模块。【免费下载链接】deltaAn open-source storage framework that enables building a Lakehouse architecture with compute engines including Spark, PrestoDB, Flink, Trino, and Hive and APIs项目地址: https://gitcode.com/GitHub_Trending/del/delta创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表