ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Apache Spark Variant 类型全解析:Parquet Variant 编码与 Shredding 规范的实现现状与限制

Apache Spark Variant 类型全解析:Parquet Variant 编码与 Shredding 规范的实现现状与限制 Apache Spark Variant 类型全解析Parquet Variant 编码与 Shredding 规范的实现现状与限制【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark本文聚焦 Apache Spark 对 Apache Parquet 项目中尚未定稿的 Variant 规范的实现Spark 将 JSON 等半结构化数据编码为value metadata双二进制表示的 Variant 值并通过 shredding 机制将 Variant 拆分为类型化列以提升存储与查询性能。读完本文你将掌握 Variant 二进制格式的编码细节、common/variant模块的核心类职责、SQL 层parse_json、variant_get等函数与相关配置项的用法以及当前实现的功能边界不支持 UUID、Time、纳秒精度 Timestamp且 shredded writes 尚无公开 API。一、Variant 是什么规范来源与 Spark 的定位Variant 是一种用于表示半结构化数据主要是 JSON的类型能够承载任意嵌套的层级结构。其二进制格式规范由Apache Parquet 项目制定规范文档分别为VariantEncoding.md与VariantShredding.md截至本仓库当前版本该规范尚未最终定稿。Spark 当前实现的是上述规范在特定提交版本上的内容并在此基础上完成了引擎内部的全链路集成。这一点在 common/variant/README.md 中有明确声明也是理解整个 Variant 功能的前提实现基于一个仍在演进中的外部规范因此后续版本存在随规范变更而调整的可能性。1.1 Spark 对规范实现的边界根据 README 的官方说明Spark 当前的实现范围与限制为已实现Variant 二进制编码Encoding与 Variant 拆分Shredding规范对应版本的内容未实现限制一尚不支持包含UUID、Time、纳秒精度 Timestamp的 Variant 值未实现限制二尚无公开 API 以启用 shredded writes拆分写入。需要说明的是从源码结构看二进制编解码层已经预留了相关能力例如 VariantUtil.java 中定义了UUIDtype info 20常量并提供getUuid/appendUuidTIMESTAMP采用微秒精度而非纳秒但 README 所述限制针对的是端到端完整支持二者并不矛盾——当前实现对 UUID 等类型的支持尚不构成完整的对外能力。1.2 模块定位与构建信息Variant 的底层实现独立于 Spark 核心位于common/variant模块构建产物为spark-variant_2.13由 common/variant/pom.xml 定义其父工程为spark-parent_2.13当前版本 5.0.0-SNAPSHOT。模块仅依赖spark-tags、spark-common-utils与jackson-coreJSON 解析/序列化设计上刻意保持轻量避免对 Spark 主体形成依赖。二、Variant 的二进制表示value 与 metadata 双二进制设计Variant 值由两个二进制数组构成value编码后的值本体与metadata字符串字典与版本信息。设计要点是读取对象/数组中的子 Variant 时通过pos偏移量直接引用同一份 value 二进制切片避免频繁复制——这一设计在 Variant.java 的构造函数注释中有明确说明value并非整段被使用而是从pos起始、长度为valueSize(value, pos)。2.1 value 的头部字节编码每个 Variant 值的第一个字节是头部字节分为两部分见 VariantUtil.java 中的常量定义高 6 位type info对于原始类型表示具体类型编号对于短字符串直接表示字符串长度MAX_SHORT_STR_SIZE 0x3F即最多 63 字节低 2 位basic type表示大类。Basic type 取值basic type值含义PRIMITIVE0原始值具体类型由 type info 决定SHORT_STR1短字符串≤63 字节内容紧跟头部字节OBJECT2对象size 字段 id 列表 字段偏移列表 字段数据ARRAY3数组size 元素偏移列表 元素数据PRIMITIVE 的 type info 取值即 Variant 支持的标量类型type info常量内容格式0NULL空1 / 2TRUE/FALSE空3 / 4 / 5 / 6INT1/INT2/INT4/INT81/2/4/8 字节小端有符号整数7DOUBLE8 字节 IEEE double8 / 9 / 10DECIMAL4/DECIMAL8/DECIMAL161 字节 scale 4/8/16 字节小端有符号整数精度上限分别为 9/18/3811DATE4 字节小端有符号整数表示距 Unix 纪元的天数12TIMESTAMP8 字节小端有符号整数表示距 Unix 纪元UTC的微秒数展示时按本地时区转换13TIMESTAMP_NTZ与 TIMESTAMP 相同的字节内容但始终按 UTC 解释14FLOAT4 字节 IEEE float15BINARY4 字节长度 内容16LONG_STR4 字节长度 内容长字符串20UUID16 字节大端源码已定义但端到端支持尚未开放值得注意的是Variant 没有独立的 Time 类型Timestamp 也只支持微秒精度这与 README 中不支持 Time 与纳秒精度 Timestamp的限制完全对应。2.2 metadata版本号与字符串字典metadata 用于压缩对象字段名等重复字符串格式为Version1 字节当前唯一允许值为 1VERSION 1掩码VERSION_MASK 0x0FDictionary size字典中字符串个数Offsets(size 1)个偏移量offsets[i]表示第 i 个字符串的起始位置字符串连续存放长度由offsets[i1] - offsets[i]推导UTF-8 字符串数据字典内容本体。metadata 头部的高 2 位还编码了偏移列表每个元素占用的字节数。对象的字段按 key 在字典中的 id 引用从而避免在每个字段处重复存储字段名字符串。2.3 大小限制与异常体系单个 Variant 的 value 与 metadata 均不得超过128 MiBSIZE_LIMIT为了测试稳定性测试环境下该上限收紧为16 MiBVariantUtil.SIZE_LIMIT通过JavaUtils.isTesting()区分配套异常包括VariantSizeLimitException超出大小上限、VariantPathTypeMismatchException路径与容器类型不匹配、以及 SQL 错误MALFORMED_VARIANT、VARIANT_CONSTRUCTOR_SIZE_LIMIT、UNKNOWN_PRIMITIVE_TYPE_IN_VARIANT等。三、核心类构建、访问、校验与 JSON 互转common/variant模块共 9 个源文件src/main/java/org/apache/spark/types/variant/职责划分清晰类核心职责Variant不可变值对象持 value/metadata/pos提供getBoolean、getLong、getDecimal、getString、getFieldByKey、getElementAtIndex、arraySize、objectSize、toJson等访问与 JSON 输出能力VariantBuilder由 JSON 解析构建 Variant实现按路径增删改、stripNulls等操作负责对象字段排序、字典构建与对象/数组头部的回填VariantUtil编码常量、类型判断getType、大小计算valueSize、标量读取、对象/数组遍历辅助handleObject/handleArray、结构校验isValidVariantVariantSchema描述 shredding schemavalue / typed_value / metadata 三元组可递归VariantShreddingWriter将 Variant 按 schema 拆分为类型化组件castShreddedShreddingUtils从拆分后的组件按规范算法重建 Variantrebuild3.1 JSON 解析与构建流程VariantBuilder.parseJson使用 Jackson 解析 JSON采用单次扫描 字段收集策略buildJson解析对象时先将每个字段的(key, id, offset)收集为FieldEntry全部解析完成后调用finishWritingObject按 key 排序、回填对象头部size、id 列表、偏移列表并用System.arraycopy将已写入的字段数据整体右移为头部腾出空间整数以最小所需宽度编码appendLong按值域自动选择 INT1/INT2/INT4/INT8见VariantBuilder.appendLong纯十进制格式且精度 ≤38 的 JSON 数字会被解析为 DECIMAL否则回退为 DOUBLEtryParseDecimal。对象字段要求按字母序排列、同一对象内不允许重复字段名解析 JSON 字符串时默认会校验 UTF-16 代理对完整性RFC 8259 §7拒绝未配对的代理项见checkValidUnicodeString以避免 Jackson 静默替换为 UFFFD 造成数据损坏。3.2 基于路径的 Variant 操作VariantBuilder同时是 SQL 中variant_delete、variant_insert、variant_set、variant_strip_nulls等函数的底层实现载体提供了一组**不可变式返回新 Variant**的路径操作deleteAtPath(v, segments)按路径删除字段/元素路径不匹配时返回语义等价的新 VariantinsertAtPath(v, segments, val)对象叶子新增字段key 已存在则抛VARIANT_DUPLICATE_KEY数组叶子在指定索引插入并右移元素越界以 null 填充setAtPath(v, segments, val, createIfMissing)替换已有值createIfMissingtrue时自动创建缺失的中间路径arrayAppendAtPath(v, segments, val)向数组末尾追加元素stripNulls(v, includeArrays)递归移除值为 null 的字段可选移除数组中的 null 元素空容器保留为{}/[]。路径段PathSegment分为ObjectKeySegment对象键与ArrayIndexSegment数组下标两类当段类型与容器类型不匹配时抛出VariantPathTypeMismatchException。所有操作都会重建 metadata因此即使没有实际删除任何内容二进制表示也可能发生变化。3.3 结构校验VariantUtil.isValidVariant(value, metadata)提供递归结构校验校验 metadata 版本、对象/数组的边界与类型信息、标量读取的合法性。其实现与toJson遍历结构一致但不强制Variant构造函数中的SIZE_LIMIT检查。它对应 SQL 层的is_valid_variant函数。四、Shredding把 Variant 拆成类型化列4.1 为什么需要 shreddingVariant 二进制是自包含的整块编码直接在列式存储中落盘会导致无法利用 Parquet 的 min/max 统计进行谓词下推、压缩率低、无法按字段级裁剪。**Shredding拆分**将 Variant 中符合 schema 的字段提取为类型化列typed_value其余字段保留在valueuntyped列中并共享同一份 metadata。Spark 对 Parquet 中 Variant 列的读写正是基于该机制。4.2 Schema 描述VariantSchemaVariantSchema.java 描述了一个合法 shredding schema核心是value / typed_value / metadata 三个字段的索引variantIdx/typedIdx/topLevelMetadataIdxtyped_value若为数组或结构体则递归包含其自身的 shredding schema元素 schema / 字段 schemametadata字段只出现在顶层递归层不包含当topLevelMetadataIdx 0 variantIdx 0 typedIdx 0时isUnshredded()返回 true表示未拆分标量 schema 支持StringType、IntegralTypeBYTE/SHORT/INT/LONG、FloatType、DoubleType、BooleanType、BinaryType、DecimalTypeprecision/scale、DateType、TimestampType、TimestampNTZType、UuidType。4.3 拆分VariantShreddingWriter.castShreddedVariantShreddingWriter.java 的castShredded按 schema 将输入 Variant 拆为ShreddedResult对象对每个字段命中 schema 的字段递归拆分进typed_value未命中的字段通过shallowAppendVariant浅拷贝进 untyped value关键正确性点浅拷贝依赖 metadata id 不变因此必须复用原 metadata缺失的 schema 字段以全字段为 null的空结果填充若 Variant 含重复字段导致同一 schema 字段被写入两次则抛MALFORMED_VARIANT数组每个元素递归拆分元素 schema 始终是含 untyped/typed 的结构体标量tryTypedShred尝试将 Variant 标量转换为目标类型——整数按目标宽度检查溢出decimal 要求精度/scale 匹配或在allowNumericScaleChanges()为 true 时允许无损的数值等价转换如整数拆成 decimal、scale 变化但不丢精度转换失败返回 null 并落入 untyped。4.4 重建ShreddingUtils.rebuild读取拆分数据时ShreddingUtils.java 按规范中的重建算法把类型化列与 untyped 残差合并回完整 Variant优先取typed_value非空时否则取value两者皆空视为输入非法MALFORMED_VARIANT。重建过程同样通过VariantBuilder完成并会拒绝untyped value 中包含已被 schema 拆分字段的数据防止重复字段。ShreddedRow接口在语义上等价于 Spark 的SpecializedGetters但被刻意独立定义避免模块对 Spark 产生依赖。五、SQL 层集成类型、函数与配置5.1VariantType与读取路径SQL 数据类型VariantType定义于 sql/api/src/main/scala/org/apache/spark/sql/types/VariantType.scala自4.0.0引入并标注UnstableAPI 尚不稳定。它是AtomicType的子类值恒可空查询规划时使用的默认大小defaultSize为 2048。JSON 数据源支持将整条记录读为单个 Variant 列singleVariantColumn选项见 docs/sql-data-sources-json.md 与 docs/sql-data-sources-csv.mdJSON 还提供prefersSingleVariantColumn将整个文件视为带顶层数组字段的文档逐元素读为 Variant 行与singleVariantColumn互斥。CSV 侧另有variantRespectInferSchema选项控制 Variant 内标量的类型推断行为。5.2 Variant 函数族variant 表达式实现集中在 variantExpressions.scala约 2100 行按引入版本分布如下函数引入版本说明parse_json/try_parse_json4.0.0字符串解析为 Varianttry_变体失败返回 nullis_variant_null4.0.0判断是否为 variant null区别于 SQL NULLto_variant_object4.0.0将 struct/array/map 转换为 Variantmap 仅限字符串键variant_get/try_variant_get4.0.0按 JSONPath 提取子值并强转为目标类型schema_of_variant/schema_of_variant_agg4.0.0 / 4.2.0推导 Variant 的 schema聚合版用于列级统计variant_from_arrays/variant_from_entries4.4.0由 keys/values 数组或 entries 结构体构造 Variant 对象variant_delete4.3.0按路径删除字段/元素variant_insert/try_variant_insert4.3.0按路径插入variant_set/try_variant_set4.3.0按路径设置值variant_strip_nulls4.3.0递归移除 null 值字段variant_explode4.3.0生成器将 Variant 数组/对象展开为行is_valid_variant4.3.0校验 Variant 二进制结构合法性variant_get的路径语法与 JSONPath 一致$.a.b[0]支持.key、[key]、[key]、[index]形式解析器为VariantPathParser基于 Scala 组合子见同文件。其底层实现VariantExpressionEvalUtils直接调用common/variant模块的Variant/VariantUtil完成二进制访问。5.3 相关配置项Variant 相关配置全部定义于 sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala均为 internal 配置供引擎内部使用不建议应用层修改。完整清单如下配置键版本默认值作用spark.sql.variant.allowDuplicateKeys4.0.0false解析 JSON 时是否允许重复键为 true 时保留同键最后出现的值spark.sql.variant.validateUnicodeInJsonParsing4.3.0true解析时拒绝含未配对 UTF-16 代理项的 JSON 字符串RFC 8259 §7为 false 恢复旧行为静默替换为 UFFFDspark.sql.variant.allowReadingShredded4.0.0trueParquet 读取时是否允许读取 shredded Variantfalse 时仅读取 unshreddedspark.sql.variant.pushVariantIntoScan4.0.0true将扫描 schema 中的 Variant 类型替换为仅含请求字段的 struct实现字段裁剪spark.sql.variant.pushVariantIntoScan.pullOutExtractions4.3.0true将字段下推扩展到聚合、join 条件、排序键、join 之上投影中的提取表达式spark.sql.variant.pushVariantIntoScan.deferCastError4.3.0false下推的严格类型转换以每行伴随错误列方式延迟抛错保持原始错误时机spark.sql.variant.writeShredding.enabled4.0.0trueParquet 写入时是否允许写 shredded Variantspark.sql.variant.shredding.maxSchemaWidth4.1.0300推断 Variant 拆分 schema 时最多创建的拆分字段数spark.sql.variant.shredding.maxSchemaDepth4.1.050推断拆分 schema 的最大遍历深度超过后按单个二进制切分spark.sql.variant.inferShreddingSchema4.1.0true写 Parquet 表时是否推断拆分 schemaspark.sql.variant.shreddedPredicatePushdown.enabled4.4.0true将 shredded 字段上的比较谓词如variant_get(v,$.a,bigint) 999下推为物理typed_value叶列上的谓词实现行组跳过spark.sql.parquet.annotateVariantLogicalType/spark.sql.parquet.ignoreVariantAnnotation——Parquet 逻辑类型标注相关其中shreddedPredicatePushdown.enabled的收益高度依赖数据布局当数据按过滤字段排序行组覆盖窄值区间且文件含多个行组时收益最大无序数据或单行组文件收益有限。该配置为纯物理扫描优化NOT_APPLICABLE绑定策略不影响查询结果。5.4 与 README无公开 API的关系README 明确声明尚无公开 API 以启用 shredded writes。这与源码现状是一致的shredding 的推断、写入与下推路径在引擎内部已经存在例如InferVariantShreddingSchema、ParquetOutputWriterWithVariantShredding、PushVariantIntoScan、PullOutVariantExtractions等实现位于 sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/ 及其 parquet 子目录但这些能力由引擎根据配置与内部推断自动触发并未向用户暴露手动指定拆分 schema 的公开 API——配置项spark.sql.variant.forceShreddingSchemaForTest的名称与注释FOR INTERNAL TESTING ONLY也印证了这一点。六、测试验证仓库为 Variant 提供了多层测试覆盖可作为理解行为的参考模块级common/variant/src/test/scala/org/apache/spark/types/variant/VariantUtf8DecodeSuite.scala 验证 UTF-8 解码表达式与类型层VariantExpressionSuite、VariantExpressionEvalUtilsSuite位于sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/variant/覆盖parse_json/variant_get等表达式语义端到端层sql/core/src/test/scala/org/apache/spark/sql/VariantEndToEndSuite.scala 覆盖parse_json/to_json往返、codegen 支持、重复键处理、schema_of_variant、cast 等场景sql/core/src/test/scala/org/apache/spark/sql/VariantSuite.scala 覆盖variant_get、variant_delete、variant_insert、try_variant_insert等路径操作含字面量与动态路径shredding 层VariantShreddingSuite、ParquetVariantShreddingSuite、VariantInferShreddingSuite、VariantShreddingFilterPushdownSuite、VariantWriteShreddingSuite位于sql/core/src/test/scala/org/apache/spark/sql/及execution/datasources/parquet/子目录验证拆分写入、schema 推断与谓词下推的正确性另有针对 XML 数据源 Variant 列行为的XmlVariantSuite。测试还验证了默认拒绝未配对 UTF-16 代理项SPARK-56654以及Variant 规范不支持 interval 类型SPARK-49985等边界行为。七、使用示例与限制总结7.1 快速上手以 JSON 数据源 SQL 函数组合的典型用法如下Spark SQL-- 将 JSON 文件整行读为单列 Variant CREATE TEMPORARY VIEW logs USING json OPTIONS (path logs.json, singleVariantColumn v); -- 解析 JSON 字符串为 Variant SELECT parse_json({a: 1, b: [true, spark]}); -- 按 JSONPath 提取并强转 SELECT variant_get(parse_json({a: {b: 42}}), $.a.b, bigint); -- 42 -- 路径操作不可变返回新 Variant SELECT variant_set(parse_json({a:1}), $.b, parse_json(2), true); -- {a:1,b:2} -- 检查 variant null 与 SQL NULL 的区别 SELECT is_variant_null(parse_json(null)); -- true SELECT is_variant_null(null); -- false注意上述函数与配置基于当前仓库源码5.0.0-SNAPSHOT功能自 4.0.0 起逐步引入若使用其他 Spark 版本请以该版本的官方文档为准。7.2 限制与注意事项规范未定稿底层格式由 Parquet 项目的 Variant 规范定义Spark 实现跟随特定提交版本规范变更可能导致格式演进类型支持不完整不支持含 UUID、Time、纳秒精度 Timestamp 的 Variant 值Timestamp 精度为微秒shredded writes 无公开 API拆分写入由引擎内部自动完成受相关配置控制用户无法手动指定拆分 schema大小上限value 与 metadata 各不超过 128 MiB测试环境 16 MiBAPI 不稳定VariantType与 variant 函数族仍标注为Unstable/实验性接口可能随版本调整。7.3 深入阅读指引规范实现与二进制格式细节common/variant/src/main/java/org/apache/spark/types/variant/VariantUtil.java、Variant.java、VariantBuilder.javaShredding 机制VariantSchema.java、VariantShreddingWriter.java、ShreddingUtils.javaSQL 函数实现variantExpressions.scala、VariantType.scala全部配置项SQLConf.scala数据源选项docs/sql-data-sources-json.md、docs/sql-data-sources-csv.md。八、结语Spark 的 Variant 支持是一套以 Apache Parquet 未定稿规范为基准、在引擎内完成全链路落地的能力二进制层通过 value/metadata 双数组与紧凑头部编码实现高密度存储shredding 层将 Variant 拆分为类型化列以换取列式存储的裁剪与谓词下推收益SQL 层则以VariantType与十余个函数构成完整的半结构化数据处理接口。理解其编码细节与当前功能边界有助于在真实业务中正确选用 Variant 能力并为后续规范演进后的迁移做好预期。【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表