Spark 4.x Variant类型详解:从JSON半结构化数据到高效存储 📅 发布时间:2026/8/27 9:12:50 👁 浏览次数: 1. 背景与核心概念1.1 半结构化数据处理到底卡在哪在日常数据开发中我们经常会遇到一类让人头疼的数据用户行为埋点、接口返回日志、第三方回调报文。这些数据有一个共同特点就是字段不固定、嵌套层级深、类型经常变化。比如一条订单日志可能今天多了个coupon字段明天某个用户没有address后天items里突然出现了null。如果按照传统二维表的方式去建模改动会非常频繁。以前面对这种数据最常见的两种方案是把整个 JSON 当作字符串保存。提前定义好 Struct 类型把 JSON 解析成一张结构化的表。第一种方案实现简单但查询很不方便。每次要读取嵌套字段都要先调用get_json_object或from_json做解析。而且 JSON 字符串本身有大量冗余的引号、冒号、空格存到 Parquet 或内存里都偏大。第二种方案性能更好但要求上游数据的 schema 非常稳定。一旦新增字段、字段类型变化或者某些嵌套层级的结构不统一ETL 就要跟着改维护成本很高。于是Spark 需要一种“既能像 JSON 一样灵活又能像结构化类型一样高效”的数据类型。这就是本文要讲的 Variant 类型出现的背景。1.2 什么是 Spark 的 Variant 类型Variant 是 Spark SQL 内置的一种动态数据类型官方从 Spark 3.3 开始以实验性功能引入后面几个版本持续补充语法、函数和存储能力到 Spark 4.x 已经越来越接近生产可用的状态。本文就以 Spark 4.x 为背景来聊一聊这个类型到底能解决什么问题以及实际用起来效果怎么样。用一句通俗的话说Variant 是一种“把 JSON 直接从字符串变成 SQL 引擎内部能够高效处理的半结构化值”的类型。它不是一种新的文件格式而是 Spark SQL 内部的类型系统扩展。一个 Variant 列可以存储对象、数组、字符串、数字、布尔值、null等不同结构的数据并且字段可以动态变化。更准确一点说Variant 采用了一种紧凑的二进制编码格式来保存半结构化数据。这种设计让它在读写、过滤、投影时比反复解析 JSON 字符串要快同时比固定 Struct 结构更灵活。类似的能力在 ClickHouse、BigQuery、DuckDB 等系统中也有对应实现只是命名和 API 不同。Spark 的 Variant 设计目标是成为处理 JSON 类数据时的默认选择之一。1.3 和 String JSON、Struct、Map 有什么区别为了理解 Variant 的价值我们把它和另外三种常见方案放在一起对比。维度STRING 存 JSONSTRUCTMAPVARIANTSchema 约束无强约束字段必须明确键和值类型固定动态字段可增可减嵌套结构字符串内部表达支持支持但值类型受限支持字段访问效率每次查询需要解析最高直接按列读取一般接近 Struct存储空间较大冗余符号多紧凑一般紧凑新增字段成本不需要改表需要改表结构不需要不需要对查询引擎友好度低高中中高从表格能看出Variant 并不是要完全替代 Struct。对于一个字段极其稳定、类型永远不变的宽表Struct 仍然是性能最好的选择。Variant 更适合的是那些“我知道大概结构但又不完全确定”的数据比如日志、事件、配置、爬虫数据。它把“先定结构再存数据”的强 schema 模式变成了“先存数据、查询时再按需解析”的灵活模式。2. 环境准备与版本说明2.1 版本选择建议使用 Variant 类型对 Spark 版本有一定要求。最早的支持从 Spark 3.3 开始但如果要用于生产环境我更建议直接使用 Spark 4.x。原因很简单Variant 从实验性功能走向成熟需要多个版本迭代4.x 在语法兼容、函数支持、Parquet 集成方面都会更完善。需要特别说明的是不同小版本的 API 和函数行为可能略有差异。本文的示例代码以 Spark 4.x 的常见行为为准。如果你使用的是 3.x 或 4.x 的某个具体 patch 版本请以官方 release notes 和对应版本文档为最终依据。版本差异不是本文能完全覆盖的内容但核心思路是一样的。如果你当前的环境还是 Spark 3.2 或更早版本那么parse_json、to_variant、variant_get这些函数很可能不存在需要先升级 Spark。2.2 开发环境Variant 相关操作可以通过 Spark SQL、PySpark、Scala 三种方式使用。考虑到大多数数据开发同学日常使用 PySpark本文核心示例用 PySpark Spark SQL 混合演示。本地开发时建议准备JDK 17 或更高版本Spark 4.x 对 JDK 版本要求比旧版更高。Python 3.9 或更高版本。PySpark直接用 pip 安装。如果不想用 Python也可以准备一个能执行 spark-sql 的 Spark 发行包。PySpark 安装命令很简单# 建议根据官方 PyPI 选择你需要的版本 pip install pyspark如果你的服务端已经有一套 Spark 环境那么直接把提交作业的方式换成 spark-submit 即可。本文示例在本地模式也能跑通。2.3 示例数据设计为了让后面的实战更有代入感我们构造一个用户行为日志的数据集。每行日志包含id日志唯一编号。raw_json原始 JSON 字符串结构包括用户信息、商品列表等。JSON 数据模拟了真实场景中的字段缺失和类型差异。这些数据会用来演示如何把字符串转成 Variant、如何提取嵌套字段、如何过滤、如何写回 Parquet以及如何观察效果。3. Variant 核心原理与基础语法3.1 用 parse_json 把字符串变成 Variant先看最简单的场景。我们有一列 JSON 字符串希望它变成真正的 Variant 类型。Spark SQL 中可以使用parse_json函数SELECT parse_json({a: 1, b: [1, 2, 3]}) AS v;执行后v的类型就是VARIANT。注意parse_json对输入字符串要求是合法 JSON。如果字符串不是合法的 JSON它会直接报错。在生产环境建议先对数据做检查或者使用容错思路在异常数据上先过滤再转换。在 PySpark 中对应的是functions.parse_jsonfrom pyspark.sql import SparkSession from pyspark.sql import functions as F spark SparkSession.builder.appName(variant_demo).getOrCreate() df spark.sql(SELECT parse_json({\a\: 1, \b\: [1, 2, 3]}) AS v) df.printSchema()输出结果中v的类型会显示为variant。这一步是后续所有 Variant 操作的基础。3.2 用 to_variant 从 Struct 或 Map 转换parse_json处理的是字符串而to_variant可以把已经存在的 Struct、Map 等类型转换成 Variant。这在从旧表迁移到新型半结构化数据时非常有用。SELECT to_variant(struct(Alice AS name, 30 AS age)) AS v;这条语句把struct(name, age)转换成了 Variant 类型。PySpark 中也可以直接使用F.to_variantfrom pyspark.sql import functions as F struct_df spark.sql(SELECT struct(Alice AS name, 30 AS age) AS s) variant_df struct_df.select(F.to_variant(F.col(s)).alias(v)) variant_df.show(truncateFalse)实际项目中我们经常把已经清洗好的 Struct 结果再转成 Variant以便统一存储逻辑。但要注意频繁的 Struct 与 Variant 互转会产生额外开销迁移前最好想清楚最终的存储类型。3.3 Variant 的字段访问语法Variant 之所以好用是因为它提供了一套嵌套字段访问语法。在 Spark SQL 中可以通过:来读取 Variant 中的子字段读取结果仍然是 Variant 类型。SELECT parse_json({user: {name: Alice}}):user.name AS user_name;这里的:user.name表示先取user字段再取name字段。返回结果仍然是 Variant因此打印出来可能是Alice这样的字符串形式但它还不是强类型的字符串列。如果是数组字段也可以使用下标SELECT parse_json({items: [{sku: A}, {sku: B}]}):items[0].sku AS first_sku;这里返回的是A类型仍然是 Variant。需要注意的是:访问语法在不同 Spark 版本中的支持程度略有差异某些版本也支持field的点号访问方式。如果遇到语法解析错误优先检查官方文档中关于 Variant 访问语法的描述或者直接改用variant_get函数。3.4 用 variant_get 提取强类型字段直接访问字段得到的是 Variant但在筛选、聚合、写入目标表时我们通常需要强类型数据。variant_get就是用来做这件事的。SELECT variant_get(v, $.user.name, STRING) AS user_name, variant_get(v, $.user.age, INT) AS user_age FROM ( SELECT parse_json({user: {name: Alice, age: 30}}) AS v ) t;variant_get接收三个参数第一个参数是 Variant 表达式。第二个参数是 JSON Path例如$.user.name。第三个参数是目标类型例如STRING、INT、DOUBLE、BOOLEAN等。如果字段不存在或者类型转换失败variant_get默认可能返回NULL具体行为取决于版本。这里要特别注意Path 中如果包含特殊字符需要按照 JSON Path 的转义规则处理。实际开发中我习惯把variant_get封装成一系列公共 SQL 片段避免每个分析任务里重复写一长串 Path。3.5 类型转换与输出Variant 也可以转换回字符串用于数据展示、落盘或下游系统消费。最简单的方式是使用 CASTSELECT CAST(parse_json({a: 1}) AS STRING) AS json_str;输出结果是紧凑的 JSON 字符串例如{a:1}。这种转换不会保留原始字符串中多余的空格和排版但数据语义是等价的。在 PySpark 中同样可以转换from pyspark.sql import functions as F variant_df spark.sql(SELECT parse_json({\a\: 1}) AS v) variant_df.select(F.col(v).cast(string).alias(json_str)).show(truncateFalse)了解了这些基础语法之后下面我们用一个完整案例把从读取到存储的整个链路串起来。4. 完整实战案例日志数据从 JSON 到 Variant 的落地4.1 创建 SparkSession 并准备原始数据我们先用 PySpark 创建一份模拟日志数据。为了贴近真实我故意让部分记录缺少字段模拟埋点数据经常出现的不规则情况。from pyspark.sql import SparkSession from pyspark.sql import functions as F spark SparkSession.builder.appName(spark_variant_demo).getOrCreate() data [ (1, {user: {name: Alice, age: 30}, items: [{sku: A, price: 9.9}]}), (2, {user: {name: Bob, age: 25}, items: [{sku: B, price: 19.9}, {sku: C, price: 3.5}]}), (3, {user: {name: Cathy}, items: []}), (4, {user: {name: David, age: 35}, items: [{sku: D, price: 29.9}]}), ] raw_df spark.createDataFrame(data, [id, raw_json]) raw_df.show(truncateFalse)执行后可以看到四行数据。其中第三条记录缺少age字段这是很典型的半结构化数据场景。此时raw_json的类型是字符串还没有任何解析。4.2 将字符串转成 Variant接下来我们把raw_json转换成 Variant得到一个名为v的新列。variant_df raw_df.withColumn(v, F.parse_json(F.col(raw_json))) variant_df.printSchema() variant_df.show(truncateFalse)输出中应该能够看到root |-- id: long (nullable true) |-- raw_json: string (nullable true) |-- v: variant (nullable true)这里的关键点是v: variant。这意味着 Spark 已经不再把 JSON 当作普通字符串保存而是转成了内部的二进制半结构化表示。从现在开始对v的字段访问不需要再去反复调用parse_json函数。4.3 使用 SQL 查询嵌套字段我们可以把variant_df注册成临时表然后使用 Spark SQL 进行复杂查询。variant_df.createOrReplaceTempView(logs) result_df spark.sql( SELECT id, variant_get(v, $.user.name, STRING) AS user_name, variant_get(v, $.user.age, INT) AS user_age, v:user.name AS user_name_variant, v:items[0] AS first_item_variant, CAST(v AS STRING) AS v_json FROM logs ) result_df.show(truncateFalse)预期结果类似--------------------------------------------------------------------------------------------------------------------- |id |user_name|user_age|user_name_variant|first_item_variant|v_json | --------------------------------------------------------------------------------------------------------------------- |1 |Alice |30 |Alice |{sku:A,price:9.9}|{user:{name:Alice,age:30},items:[{sku:A,price:9.9}]}| |2 |Bob |25 |Bob |{sku:B,price:19.9}|{user:{name:Bob,age:25},items:[{sku:B,price:19.9},{sku:C,price:3.5}]}| |3 |Cathy |null |Cathy |null |{user:{name:Cathy},items:[]} | |4 |David |35 |David |{sku:D,price:29.9}|{user:{name:David,age:35},items:[{sku:D,price:29.9}]}| ---------------------------------------------------------------------------------------------------------------------可以看到variant_get返回的是强类型字段user_age直接是整数可用于后续数值运算。v:user.name返回的是 Variant 类型的字段输出时会表现为带引号的字符串。v:items[0]返回了数组第一个元素本身仍然是一个 Variant 对象。缺失的age字段在user_age中显示为null不会导致任务失败。4.4 基于 Variant 做过滤与统计Variant 虽然灵活但一旦转出强类型字段就可以正常参与过滤、聚合、排序等操作。例如我们要筛选出年龄大于等于 28 岁的用户并且按年龄倒序排列SELECT id, variant_get(v, $.user.name, STRING) AS user_name, variant_get(v, $.user.age, INT) AS user_age FROM logs WHERE variant_get(v, $.user.age, INT) IS NOT NULL AND variant_get(v, $.user.age, INT) 28 ORDER BY user_age DESC执行结果会返回 id 为 4 和 1 的两条记录。这里有一点需要提醒如果直接对user_age做比较但 Variant 值是字符串类型可能会出现类型转换问题。因此在过滤前先通过variant_get转成INT是更稳妥的做法。再比如我们想统计每个用户的商品数量。由于 Variant 内部是动态类型直接对v:items求长度可能不是所有版本都支持。更稳妥的方案是先把它提取成字符串再交给 JSON 函数处理或者通过variant_get转成目标数组类型。实际项目中我推荐把这类常用字段提取逻辑封装成一个视图或公共表让分析人员直接使用已经转好的强类型字段而不是每次都写一长串variant_get。4.5 写入 Parquet 并重新读取Variant 的一个重要价值是它能够被存储到列式文件中并且在下次读取时保持类型不变。下面我们把variant_df写入 Parquet。variant_df.write.mode(overwrite).format(parquet).save(./variant_demo_parquet)重新读取parquet_df spark.read.parquet(./variant_demo_parquet) parquet_df.printSchema() parquet_df.select(id, F.col(v).cast(string).alias(v_json)).show(truncateFalse)读取后v列仍然是variant类型不需要再做一次parse_json。这意味着只要在写入时完成一次解析后续每次读取都省去了重复解析 JSON 的开销。这一点对高频查询场景非常重要。4.6 如何观察“效果咋样”很多同学会问Variant 到底比 String 快多少、省多少空间这个问题很难给一个统一的数字因为效果取决于数据内容、压缩算法、访问模式。但我们可以通过一个简单的对比实验来感受差异。先把 JSON 字符串原样写入一个 Parquetstring_df raw_df.select(id, F.col(raw_json).alias(v)) string_df.write.mode(overwrite).format(parquet).save(./string_json_parquet)然后使用系统命令查看两个目录的大小du -sh ./string_json_parquet ./variant_demo_parquet可以预期Variant 存储的 Parquet 在大多数情况下比纯 JSON 字符串要小因为去掉了很多 JSON 语法符号并且内部编码更紧凑。但由于 Parquet 自带压缩在小样本数据上差异可能不明显。数据量越大、JSON 重复结构越多差异通常越明显。同样也可以用查询计划来观察result_df.explain()通过explain查看执行计划可以看到是否还有额外的 JSON 解析步骤。如果数据已经以 Variant 形式存储查询计划应该更简洁。5. 常见问题与排查思路5.1 parse_json 遇到非法 JSON问题现象常见原因解决思路作业报错提示 JSON 解析失败源数据里混入了非 JSON 字符串先清洗数据或使用容错函数再执行转换parse_json是严格解析。如果上游数据中偶尔出现空字符串、JSON 截断、中文注释等作业就会失败。推荐先把明显非法的数据过滤掉或者使用try_parse_json这类容错接口。如果版本里没有容错函数可以先用get_json_object或from_json做一次探测把非法数据排除后再转。5.2 Variant 字段访问语法报错问题现象常见原因解决思路v:user.name语法不生效版本不支持该语法或者字段名需要转义改用variant_get并确认 Path 写法不同 Spark 小版本对:和.的支持程度不同。遇到语法报错时先检查官方文档再看是否要把variant_get作为统一方案。字段名如果包含特殊字符建议用variant_get(v, $.字段名, STRING)处理避免写出无法解析的访问表达式。5.3 下游系统不识别 Parquet 中的 Variant问题现象常见原因解决思路其他引擎读取 Parquet 时Variant 列变成不可读类型下游引擎不支持 Variant 物理类型落表时额外提供一列 JSON 字符串或明确下游版本支持范围Variant 写入 Parquet 后依赖 Parquet 对 Variant 类型的支持。如果下游是 Presto、Hive 或其他旧引擎可能无法直接读取这种列类型。生产环境里如果存在多引擎消费我会建议在表中同时保留v和v_json两列让不支持 Variant 的引擎也能通过字符串列读取。5.4 自定义 UDF 处理 Variant 失败问题现象常见原因解决思路UDF 接收 Variant 后无法正常处理自定义函数不支持动态类型在 UDF 外先 CAST 成字符串或 StructSpark 原生函数对 Variant 的支持正在逐步完善但自定义 UDF 通常没有自动适配 Variant 的逻辑。遇到这种场景最稳妥的方式是在调用 UDF 前先把 Variant 列转成字符串或强类型 Struct再把转换结果传入函数。5.5 性能没有明显提升问题现象常见原因解决思路使用 Variant 后查询还是很慢数据仍然是字符串查询时才 parse在写入阶段完成 parse并使用 Variant 存储Variant 的优势建立在“只解析一次”的基础上。如果表里存的仍然是 JSON 字符串每次查询都要临时转换收益就会大打折扣。正确做法是在 ETL 写入时把合法的 JSON 字符串统一转成 Variant之后查询直接使用转换后的列。6. 最佳实践与工程建议6.1 在写入阶段完成转换这是使用 Variant 最重要的一条原则。不要让下游分析人员每次查询都执行parse_json而应该在数据落地前完成解析和转换。也就是说ODS 层可以保留原始 JSON 字符串但到了 DWD 或 DWS 层如果确实需要半结构化能力就应该把核心字段转成 Variant 存储。这样做的好处有三个一是查询时节省重复解析开销二是写入 Parquet 后存储更紧凑三是下游接口可以统一读取强类型字段。6.2 强 Schema 数据还是要用 StructVariant 不等于万能类型。如果一个业务表的所有字段都固定类型也不会变化那么使用 Struct 可以获得更好的性能和更严格的约束。Variant 适合的是字段容易变化、结构不规则的数据比如埋点、日志、配置、AI 模型输出等。如果业务上已经存在一个字段长期稳定的核心宽表不要因为“Variant 听起来很强”就盲目改造。从 Struct 改成 Variant 容易但后续查询性能可能下降而且下游依赖 schema 的任务需要同步调整。6.3 为下游准备兼容视图在跨团队协作时下游消费者不一定了解 Variant。建议在表之上创建一层视图把常用的嵌套字段通过variant_get提取成强类型列。这样熟悉 Variant 的团队可以直接操作原始列。不熟悉 Variant 的团队只查询视图中的普通列。字段名统一维护避免每个任务里写不一样的 JSON Path。视图本质上不占额外存储是一种成本很低的兼容手段。6.4 关注版本和函数支持范围Variant 是相对较新的类型不同版本对函数、语法、Parquet 写入支持有差异。上线前建议做一轮版本验证重点检查parse_json、to_variant、variant_get是否可用。字段访问语法是否正常。Variant 写入 Parquet 后能否被预期版本读取。压缩算法和文件大小是否符合预期。在升级 Spark 版本时也要把 Variant 相关的查询加入回归测试避免某个函数行为变化影响线上任务。6.5 生产环境变更要谨慎如果要在生产表上新增 Variant 列或修改列类型建议按常规变更流程走。先用一份测试数据在小范围环境验证再通过变更单执行。涉及核心业务表时保留一份旧结构的备份确保回滚路径明确。不要因为“只是换个类型”就跳过验证Variant 在 Parquet 物理层的行为和普通列不一样风险往往出现在下游兼容性上。性能评估也不要只看官方 benchmark。不同数据集的 JSON 重复度、嵌套深度、压缩算法都会影响结果。最可靠的方式是拿自己业务里最典型的一批数据做一次 String 和 Variant 的对比测试记录文件大小、查询耗时、资源消耗再做决策。7. 总结与学习路线这篇内容从半结构化数据的痛点讲起介绍了 Spark 4.x 中 Variant 类型的基本概念、核心语法并通过一个完整的日志解析案例演示了从字符串 JSON 到 Variant、再到 Parquet 存储的落地流程。你现在应该已经掌握了几个关键点Variant 是 Spark SQL 内部的一种动态数据类型适合处理 JSON 这类结构不固定的数据。parse_json和to_variant负责创建 Variantvariant_get负责提取强类型字段。Variant 最推荐的使用方式是写入时完成转换查询时直接使用避免反复解析。存储到 Parquet 后Variant 类型可以继续保持这是它比 String JSON 有优势的重要原因。下一步你可以继续深入研究官方发布的 Variant 设计文档了解它的二进制编码细节。如果工作中有真实的埋点数据建议直接拿一份历史数据跑一遍本文的流程对比文件大小和查询耗时。只有结合你自己的数据特征才能真正判断出 Spark 4.x 的 Variant 在你的场景里效果怎么样。如果在跑实验的过程中遇到奇怪的报错欢迎带着版本号和报错信息在评论区交流。这类新特性的坑往往需要大家一起踩过才知道怎么避开。