如何在Spark SQL中条件式展开多数组列?
Spark SQL条件式展开多数组列解决方案
针对你提出的按索引匹配数组元素、单元素数组重复映射的需求,可以通过posexplode结合条件判断实现,具体方案如下:
核心SQL实现
SELECT col1, col2_val, CASE WHEN size(col3) = 1 THEN col3[0] ELSE col3[pos] END AS col3_val FROM my_table LATERAL VIEW posexplode(col2) exploded AS pos, col2_val
逻辑说明
posexplode的作用:和普通explode不同,它会同时返回数组元素的**索引(pos)**和对应的值,这是实现"相同索引元素映射到同一行"的关键,避免了双重explode导致的笛卡尔积问题。- 条件处理col3:
- 当
col3长度为1时,直接取第一个元素col3[0],实现单元素对所有col2元素的重复映射; - 当
col3长度大于1时,通过索引pos取对应位置的元素,保证和col2的元素按索引匹配。
- 当
执行结果
用你的示例数据运行上述SQL,会得到符合预期的输出:
+----+--------+--------+ |col1|col2_val|col3_val| +----+--------+--------+ | 123| id_1| tim| | 123| id_2| steve| | 456| id_3| jenny| | 456| id_4| jenny| +----+--------+--------+
边界情况补充
如果存在col3长度既不是1也不等于col2长度的情况,可以添加索引越界防护:
SELECT col1, col2_val, CASE WHEN size(col3) = 1 THEN col3[0] WHEN pos < size(col3) THEN col3[pos] ELSE NULL -- 可替换为你需要的默认值 END AS col3_val FROM my_table LATERAL VIEW posexplode(col2) exploded AS pos, col2_val
可选DataFrame API写法
如果你需要用DataFrame代码实现,逻辑和SQL一致:
from pyspark.sql import functions as F df = spark.createDataFrame( [(123, ["id_1", "id_2"], ["tim", "steve"]), (456, ["id_3", "id_4"], ["jenny"])], ["col1", "col2", "col3"] ) result_df = df.select( "col1", F.posexplode("col2").alias("pos", "col2_val"), "col3" ).withColumn( "col3_val", F.when(F.size("col3") == 1, F.col("col3")[0]) .otherwise(F.col("col3")[F.col("pos")]) ).drop("pos", "col3") result_df.show()
内容的提问来源于stack exchange,提问作者Dumbledore__
相关产品推荐
相关产品推荐

