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 中定义了UUID(type info = 20)常量并提供getUuid/appendUuid,TIMESTAMP采用微秒精度而非纳秒),但 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-core(JSON 解析/序列化),设计上刻意保持轻量,避免对 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 | 值 | 含义 |
|---|---|---|
PRIMITIVE | 0 | 原始值,具体类型由 type info 决定 |
SHORT_STR | 1 | 短字符串(≤63 字节),内容紧跟头部字节 |
OBJECT | 2 | 对象:size + 字段 id 列表 + 字段偏移列表 + 字段数据 |
ARRAY | 3 | 数组:size + 元素偏移列表 + 元素数据 |
PRIMITIVE 的 type info 取值(即 Variant 支持的标量类型):
| type info | 常量 | 内容格式 |
|---|---|---|
| 0 | NULL | 空 |
| 1 / 2 | TRUE/FALSE | 空 |
| 3 / 4 / 5 / 6 | INT1/INT2/INT4/INT8 | 1/2/4/8 字节小端有符号整数 |
| 7 | DOUBLE | 8 字节 IEEE double |
| 8 / 9 / 10 | DECIMAL4/DECIMAL8/DECIMAL16 | 1 字节 scale + 4/8/16 字节小端有符号整数(精度上限分别为 9/18/38) |
| 11 | DATE | 4 字节小端有符号整数,表示距 Unix 纪元的天数 |
| 12 | TIMESTAMP | 8 字节小端有符号整数,表示距 Unix 纪元(UTC)的微秒数,展示时按本地时区转换 |
| 13 | TIMESTAMP_NTZ | 与 TIMESTAMP 相同的字节内容,但始终按 UTC 解释 |
| 14 | FLOAT | 4 字节 IEEE float |
| 15 | BINARY | 4 字节长度 + 内容 |
| 16 | LONG_STR | 4 字节长度 + 内容(长字符串) |
| 20 | UUID | 16 字节大端(源码已定义,但端到端支持尚未开放) |
值得注意的是Variant 没有独立的 Time 类型,Timestamp 也只支持微秒精度,这与 README 中"不支持 Time 与纳秒精度 Timestamp"的限制完全对应。
2.2 metadata:版本号与字符串字典
metadata 用于压缩对象字段名等重复字符串,格式为:
- Version:1 字节,当前唯一允许值为 1(
VERSION = 1,掩码VERSION_MASK = 0x0F); - Dictionary size:字典中字符串个数;
- Offsets:
(size + 1)个偏移量,offsets[i]表示第 i 个字符串的起始位置,字符串连续存放,长度由offsets[i+1] - offsets[i]推导; - UTF-8 字符串数据:字典内容本体。
metadata 头部的高 2 位还编码了偏移列表每个元素占用的字节数。对象的字段按 key 在字典中的 id 引用,从而避免在每个字段处重复存储字段名字符串。
2.3 大小限制与异常体系
- 单个 Variant 的 value 与 metadata 均不得超过128 MiB(
SIZE_LIMIT);为了测试稳定性,测试环境下该上限收紧为16 MiB(VariantUtil.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)、结构校验(isValidVariant) |
VariantSchema | 描述 shredding schema(value / typed_value / metadata 三元组,可递归) |
VariantShreddingWriter | 将 Variant 按 schema 拆分为类型化组件(castShredded) |
ShreddingUtils | 从拆分后的组件按规范算法重建 Variant(rebuild) |
3.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,否则回退为 DOUBLE(
tryParseDecimal)。
对象字段要求按字母序排列、同一对象内不允许重复字段名;解析 JSON 字符串时,默认会校验 UTF-16 代理对完整性(RFC 8259 §7),拒绝未配对的代理项(见checkValidUnicodeString),以避免 Jackson 静默替换为 U+FFFD 造成数据损坏。
3.2 基于路径的 Variant 操作
VariantBuilder同时是 SQL 中variant_delete、variant_insert、variant_set、variant_strip_nulls等函数的底层实现载体,提供了一组**不可变式(返回新 Variant)**的路径操作:
deleteAtPath(v, segments):按路径删除字段/元素;路径不匹配时返回语义等价的新 Variant;insertAtPath(v, segments, val):对象叶子新增字段(key 已存在则抛VARIANT_DUPLICATE_KEY);数组叶子在指定索引插入并右移元素,越界以 null 填充;setAtPath(v, segments, val, createIfMissing):替换已有值,createIfMissing=true时自动创建缺失的中间路径;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 为什么需要 shredding
Variant 二进制是自包含的整块编码,直接在列式存储中落盘会导致:无法利用 Parquet 的 min/max 统计进行谓词下推、压缩率低、无法按字段级裁剪。**Shredding(拆分)**将 Variant 中符合 schema 的字段提取为类型化列(typed_value),其余字段保留在value(untyped)列中,并共享同一份 metadata。Spark 对 Parquet 中 Variant 列的读写正是基于该机制。
4.2 Schema 描述:VariantSchema
VariantSchema.java 描述了一个合法 shredding schema,核心是value / typed_value / metadata 三个字段的索引(variantIdx/typedIdx/topLevelMetadataIdx):
typed_value若为数组或结构体,则递归包含其自身的 shredding schema(元素 schema / 字段 schema);metadata字段只出现在顶层,递归层不包含;- 当
topLevelMetadataIdx >= 0 && variantIdx >= 0 && typedIdx < 0时,isUnshredded()返回 true,表示未拆分; - 标量 schema 支持:
StringType、IntegralType(BYTE/SHORT/INT/LONG)、FloatType、DoubleType、BooleanType、BinaryType、DecimalType(precision/scale)、DateType、TimestampType、TimestampNTZType、UuidType。
4.3 拆分:VariantShreddingWriter.castShredded
VariantShreddingWriter.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引入并标注@Unstable(API 尚不稳定)。它是AtomicType的子类,值恒可空,查询规划时使用的默认大小(defaultSize)为 2048。
JSON 数据源支持将整条记录读为单个 Variant 列:singleVariantColumn选项(见 docs/sql-data-sources-json.md 与 docs/sql-data-sources-csv.md);JSON 还提供prefersSingleVariantColumn(将整个文件视为带顶层数组字段的文档,逐元素读为 Variant 行,与singleVariantColumn互斥)。CSV 侧另有variantRespectInferSchema选项控制 Variant 内标量的类型推断行为。
5.2 Variant 函数族
variant 表达式实现集中在 variantExpressions.scala(约 2100 行),按引入版本分布如下:
| 函数 | 引入版本 | 说明 |
|---|---|---|
parse_json/try_parse_json | 4.0.0 | 字符串解析为 Variant,try_变体失败返回 null |
is_variant_null | 4.0.0 | 判断是否为 variant null(区别于 SQL NULL) |
to_variant_object | 4.0.0 | 将 struct/array/map 转换为 Variant(map 仅限字符串键) |
variant_get/try_variant_get | 4.0.0 | 按 JSONPath 提取子值并强转为目标类型 |
schema_of_variant/schema_of_variant_agg | 4.0.0 / 4.2.0 | 推导 Variant 的 schema(聚合版用于列级统计) |
variant_from_arrays/variant_from_entries | 4.4.0 | 由 keys/values 数组(或 entries 结构体)构造 Variant 对象 |
variant_delete | 4.3.0 | 按路径删除字段/元素 |
variant_insert/try_variant_insert | 4.3.0 | 按路径插入 |
variant_set/try_variant_set | 4.3.0 | 按路径设置值 |
variant_strip_nulls | 4.3.0 | 递归移除 null 值字段 |
variant_explode | 4.3.0 | 生成器,将 Variant 数组/对象展开为行 |
is_valid_variant | 4.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.allowDuplicateKeys | 4.0.0 | false | 解析 JSON 时是否允许重复键;为 true 时保留同键最后出现的值 |
spark.sql.variant.validateUnicodeInJsonParsing | 4.3.0 | true | 解析时拒绝含未配对 UTF-16 代理项的 JSON 字符串(RFC 8259 §7);为 false 恢复旧行为(静默替换为 U+FFFD) |
spark.sql.variant.allowReadingShredded | 4.0.0 | true | Parquet 读取时是否允许读取 shredded Variant;false 时仅读取 unshredded |
spark.sql.variant.pushVariantIntoScan | 4.0.0 | true | 将扫描 schema 中的 Variant 类型替换为仅含请求字段的 struct,实现字段裁剪 |
spark.sql.variant.pushVariantIntoScan.pullOutExtractions | 4.3.0 | true | 将字段下推扩展到聚合、join 条件、排序键、join 之上投影中的提取表达式 |
spark.sql.variant.pushVariantIntoScan.deferCastError | 4.3.0 | false | 下推的严格类型转换以"每行伴随错误列"方式延迟抛错,保持原始错误时机 |
spark.sql.variant.writeShredding.enabled | 4.0.0 | true | Parquet 写入时是否允许写 shredded Variant |
spark.sql.variant.shredding.maxSchemaWidth | 4.1.0 | 300 | 推断 Variant 拆分 schema 时最多创建的拆分字段数 |
spark.sql.variant.shredding.maxSchemaDepth | 4.1.0 | 50 | 推断拆分 schema 的最大遍历深度,超过后按单个二进制切分 |
spark.sql.variant.inferShreddingSchema | 4.1.0 | true | 写 Parquet 表时是否推断拆分 schema |
spark.sql.variant.shreddedPredicatePushdown.enabled | 4.4.0 | true | 将 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 MiB);
- API 不稳定:
VariantType与 variant 函数族仍标注为@Unstable/实验性,接口可能随版本调整。
7.3 深入阅读指引
- 规范实现与二进制格式细节:common/variant/src/main/java/org/apache/spark/types/variant/VariantUtil.java、Variant.java、VariantBuilder.java;
- Shredding 机制:VariantSchema.java、VariantShreddingWriter.java、ShreddingUtils.java;
- SQL 函数实现: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),仅供参考