如何在Apache Spark中按结构字段分组合并两个ArrayType列
解决Spark中合并不同顺序/集合的ArrayType列问题
你需要合并单行数据中的两个结构化ArrayType列carColor(结构{type, color, owner})和carPrice(结构{type, price}),生成新列cars,要求:
- 按
type分组匹配,忽略原数组的顺序差异 - 两个数组的车型集合可能不同,缺失字段补
null(color)或0(price) - 现有代码仅在车型顺序、数量完全一致时有效,需要适配通用场景
原始数据集示例
+--+--+------------------------------------------+-------------------------------+ |C1|C2| carColor | carPrice | +--+--+------------------------------------------+-------------------------------+ | 1|A |[{type: BMW, color: red, owner: Alice}, | [{type: Toyota, price: 30000},| | | | {type: Ford, color: blue, owner: Bob}, | {type: BMW, price: 50000}] | | | | {type: Honda, color: pink, owner: Chuck}]| | +--+--+------------------------------------------+-------------------------------+
期望输出数据集示例
+--+--+--------------------------------------------+ |C1|C2| cars | +--+--+--------------------------------------------+ | 1|A |[{type: BMW, color: red, price: 50000}, | | | | {type: Ford, color: blue, price: 0}, | | | | {type: Honda, color: pink, price: 0}, | | | | {type: Toyota, color: null, price: 30000}] | +--+--+--------------------------------------------+
解决方案
核心思路是先将两个数组拆分成单行,基于type做全关联补全缺失数据,再聚合回数组,彻底解决顺序、集合差异的问题。
方法一:Java代码实现
private static Dataset<Row> addCars(Dataset<Row> dataset) { return dataset // 展开carColor数组,提取type和color .withColumn("color_item", explode(col("carColor"))) .withColumn("type_c", col("color_item.type")) .withColumn("color", col("color_item.color")) // 展开carPrice数组,提取type和price .withColumn("price_item", explode(col("carPrice"))) .withColumn("type_p", col("price_item.type")) .withColumn("price", col("price_item.price")) // 全关联补全所有车型,统一type字段 .select( col("C1"), col("C2"), coalesce(col("type_c"), col("type_p")).alias("type"), col("color"), col("price") ) // 按主键+type分组,确保同一车型只保留一条有效数据 .groupBy("C1", "C2", "type") .agg( first(col("color"), true).alias("color"), coalesce(first(col("price"), true), lit(0)).alias("price") ) // 聚合回数组,组装目标结构 .groupBy("C1", "C2") .agg( collect_list(struct(col("type"), col("color"), col("price"))).alias("cars") ); }
方法二:Spark SQL语句实现
如果习惯用SQL语法,可注册临时表后执行:
WITH color_expanded AS ( SELECT C1, C2, color_item.type AS type, color_item.color AS color FROM your_table LATERAL VIEW explode(carColor) AS color_item ), price_expanded AS ( SELECT C1, C2, price_item.type AS type, price_item.price AS price FROM your_table LATERAL VIEW explode(carPrice) AS price_item ), joined_data AS ( SELECT COALESCE(c.C1, p.C1) AS C1, COALESCE(c.C2, p.C2) AS C2, COALESCE(c.type, p.type) AS type, c.color, p.price FROM color_expanded c FULL OUTER JOIN price_expanded p ON c.C1 = p.C1 AND c.C2 = p.C2 AND c.type = p.type ) SELECT C1, C2, COLLECT_LIST(STRUCT(type, color, COALESCE(price, 0) AS price)) AS cars FROM joined_data GROUP BY C1, C2;
关键步骤说明
- 数组展开:用
explode将ArrayType列拆分为单行,让每个车型成为独立数据行,方便后续匹配 - 全关联:通过
FULL OUTER JOIN保留两个数组的所有车型,避免遗漏任何一方存在的车型 - 缺失值补全:
color缺失时保留null,符合需求price缺失时用COALESCE替换为0
- 聚合回数组:用
collect_list将同一主键下的所有车型数据重新组装成ArrayType列
内容的提问来源于stack exchange,提问作者antemerid
相关产品推荐
相关产品推荐

