Apache Spark Variant 类型全解析:Parquet Variant 编码与 Shredding 规范的实现现状与限制
2026/9/19 10:15:28 网站建设 项目流程

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_jsonvariant_get等函数与相关配置项)的用法,以及当前实现的功能边界(不支持 UUID、Time、纳秒精度 Timestamp,且 shredded writes 尚无公开 API)。

一、Variant 是什么:规范来源与 Spark 的定位

Variant 是一种用于表示半结构化数据(主要是 JSON)的类型,能够承载任意嵌套的层级结构。其二进制格式规范由Apache Parquet 项目制定(规范文档分别为VariantEncoding.mdVariantShredding.md),截至本仓库当前版本,该规范尚未最终定稿

Spark 当前实现的是上述规范在特定提交版本上的内容,并在此基础上完成了引擎内部的全链路集成。这一点在 common/variant/README.md 中有明确声明,也是理解整个 Variant 功能的前提:实现基于一个仍在演进中的外部规范,因此后续版本存在随规范变更而调整的可能性。

1.1 Spark 对规范实现的边界

根据 README 的官方说明,Spark 当前的实现范围与限制为:

  • 已实现:Variant 二进制编码(Encoding)与 Variant 拆分(Shredding)规范对应版本的内容;
  • 未实现(限制一):尚不支持包含UUIDTime纳秒精度 Timestamp的 Variant 值;
  • 未实现(限制二):尚无公开 API 以启用 shredded writes(拆分写入)。

需要说明的是,从源码结构看,二进制编解码层已经预留了相关能力(例如 VariantUtil.java 中定义了UUID(type 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-tagsspark-common-utilsjackson-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含义
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 double
8 / 9 / 10DECIMAL4/DECIMAL8/DECIMAL161 字节 scale + 4/8/16 字节小端有符号整数(精度上限分别为 9/18/38)
11DATE4 字节小端有符号整数,表示距 Unix 纪元的天数
12TIMESTAMP8 字节小端有符号整数,表示距 Unix 纪元(UTC)的微秒数,展示时按本地时区转换
13TIMESTAMP_NTZ与 TIMESTAMP 相同的字节内容,但始终按 UTC 解释
14FLOAT4 字节 IEEE float
15BINARY4 字节长度 + 内容
16LONG_STR4 字节长度 + 内容(长字符串)
20UUID16 字节大端(源码已定义,但端到端支持尚未开放)

值得注意的是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 MiBSIZE_LIMIT);为了测试稳定性,测试环境下该上限收紧为16 MiBVariantUtil.SIZE_LIMIT,通过JavaUtils.isTesting()区分);
  • 配套异常包括:VariantSizeLimitException(超出大小上限)、VariantPathTypeMismatchException(路径与容器类型不匹配)、以及 SQL 错误MALFORMED_VARIANTVARIANT_CONSTRUCTOR_SIZE_LIMITUNKNOWN_PRIMITIVE_TYPE_IN_VARIANT等。

三、核心类:构建、访问、校验与 JSON 互转

common/variant模块共 9 个源文件(src/main/java/org/apache/spark/types/variant/),职责划分清晰:

核心职责
Variant不可变值对象:持 value/metadata/pos;提供getBooleangetLonggetDecimalgetStringgetFieldByKeygetElementAtIndexarraySizeobjectSizetoJson等访问与 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):

  1. 解析对象时先将每个字段的(key, id, offset)收集为FieldEntry,全部解析完成后调用finishWritingObject按 key 排序、回填对象头部(size、id 列表、偏移列表),并用System.arraycopy将已写入的字段数据整体右移为头部腾出空间;
  2. 整数以最小所需宽度编码(appendLong按值域自动选择 INT1/INT2/INT4/INT8,见VariantBuilder.appendLong);
  3. 纯十进制格式且精度 ≤38 的 JSON 数字会被解析为 DECIMAL,否则回退为 DOUBLE(tryParseDecimal)。

对象字段要求按字母序排列、同一对象内不允许重复字段名;解析 JSON 字符串时,默认会校验 UTF-16 代理对完整性(RFC 8259 §7),拒绝未配对的代理项(见checkValidUnicodeString),以避免 Jackson 静默替换为 U+FFFD 造成数据损坏。

3.2 基于路径的 Variant 操作

VariantBuilder同时是 SQL 中variant_deletevariant_insertvariant_setvariant_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 支持:StringTypeIntegralType(BYTE/SHORT/INT/LONG)、FloatTypeDoubleTypeBooleanTypeBinaryTypeDecimalType(precision/scale)、DateTypeTimestampTypeTimestampNTZTypeUuidType

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_json4.0.0字符串解析为 Variant,try_变体失败返回 null
is_variant_null4.0.0判断是否为 variant null(区别于 SQL NULL)
to_variant_object4.0.0将 struct/array/map 转换为 Variant(map 仅限字符串键)
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 恢复旧行为(静默替换为 U+FFFD)
spark.sql.variant.allowReadingShredded4.0.0trueParquet 读取时是否允许读取 shredded Variant;false 时仅读取 unshredded
spark.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 Variant
spark.sql.variant.shredding.maxSchemaWidth4.1.0300推断 Variant 拆分 schema 时最多创建的拆分字段数
spark.sql.variant.shredding.maxSchemaDepth4.1.050推断拆分 schema 的最大遍历深度,超过后按单个二进制切分
spark.sql.variant.inferShreddingSchema4.1.0true写 Parquet 表时是否推断拆分 schema
spark.sql.variant.shreddedPredicatePushdown.enabled4.4.0true将 shredded 字段上的比较谓词(如variant_get(v,'$.a','bigint') > 999)下推为物理typed_value叶列上的谓词,实现行组跳过
spark.sql.parquet.annotateVariantLogicalType/spark.sql.parquet.ignoreVariantAnnotationParquet 逻辑类型标注相关

其中shreddedPredicatePushdown.enabled的收益高度依赖数据布局:当数据按过滤字段排序(行组覆盖窄值区间)且文件含多个行组时收益最大;无序数据或单行组文件收益有限。该配置为纯物理扫描优化(NOT_APPLICABLE绑定策略),不影响查询结果。

5.4 与 README"无公开 API"的关系

README 明确声明"尚无公开 API 以启用 shredded writes"。这与源码现状是一致的:shredding 的推断、写入与下推路径在引擎内部已经存在(例如InferVariantShreddingSchemaParquetOutputWriterWithVariantShreddingPushVariantIntoScanPullOutVariantExtractions等实现位于 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 解码;
  • 表达式与类型层VariantExpressionSuiteVariantExpressionEvalUtilsSuite(位于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_getvariant_deletevariant_inserttry_variant_insert等路径操作(含字面量与动态路径);
  • shredding 层VariantShreddingSuiteParquetVariantShreddingSuiteVariantInferShreddingSuiteVariantShreddingFilterPushdownSuiteVariantWriteShreddingSuite(位于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 限制与注意事项

  1. 规范未定稿:底层格式由 Parquet 项目的 Variant 规范定义,Spark 实现跟随特定提交版本,规范变更可能导致格式演进;
  2. 类型支持不完整:不支持含 UUID、Time、纳秒精度 Timestamp 的 Variant 值;Timestamp 精度为微秒;
  3. shredded writes 无公开 API:拆分写入由引擎内部自动完成(受相关配置控制),用户无法手动指定拆分 schema;
  4. 大小上限:value 与 metadata 各不超过 128 MiB(测试环境 16 MiB);
  5. 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询