PySpark动态实现多Hierarchy列转单层级列(Unpivot/Reduce)
动态处理任意数量层级列的PySpark DataFrame转换
问题背景
原始PySpark DataFrame结构如下:
df = spark.createDataFrame( [ ("D1", "D2", "H1", None, None), ("D1", "D2", "H1", "H2", None), ("D1", "D2", "H1", "H2", "H3") ], ["Dimension1", "Dimention2", "Hierarchy1", "Hierarchy2", "Hierarchy3"] )
需要转换为如下结构的DataFrame:
new_df = spark.createDataFrame( [ ("D1", "D2", "H1"), ("D1", "D2", "H2"), ("D1", "D2", "H3") ], ["Dimension1", "Dimention2", "Hierarchy"] )
转换逻辑:
- 仅Hierarchy1非空时,Hierarchy列取值为"H1"
- Hierarchy1、Hierarchy2非空且Hierarchy3为空时,取值为"H2"
- 三者都非空时,取值为"H3"
固定列数的实现代码已可运行,但层级列数量不固定,需要适配任意数量层级列的通用方案。
解决方案
方法一:利用数组取最后非空元素(简洁高效)
由于层级列值正好是对应的层级名称(如Hierarchy1列的值为"H1"),可以直接将所有层级列转为数组,过滤空值后取最后一个元素:
from pyspark.sql import functions as F # 定义层级列列表,支持任意数量 hierarchies = ["Hierarchy1", "Hierarchy2", "Hierarchy3"] new_df = df.select( "Dimension1", "Dimention2", # 将层级列转为数组 → 过滤空值 → 取最后一个元素 F.element_at(F.filter(F.array(*hierarchies), lambda x: x.isNotNull()), -1).alias("Hierarchy") ) new_df.show()
方法二:动态构建条件链(通用场景)
如果层级列的值不是对应的层级名称,需要基于列是否非空的逻辑动态生成when条件链:
from pyspark.sql import functions as F from functools import reduce hierarchies = ["Hierarchy1", "Hierarchy2", "Hierarchy3"] # 为每个层级生成对应的条件:前i+1个层级非空,后续层级为空 condition_list = [] for idx, col_name in enumerate(hierarchies): # 生成"前i+1列非空"的条件 non_null_cond = reduce(lambda a, b: a & b, [F.col(c).isNotNull() for c in hierarchies[:idx+1]]) # 生成"后续列全为空"的条件(如果有后续列) if idx < len(hierarchies) - 1: null_cond = reduce(lambda a, b: a & b, [F.col(c).isNull() for c in hierarchies[idx+1:]]) full_cond = non_null_cond & null_cond else: # 最后一个层级仅需自身非空(前面的已由前面的条件覆盖) full_cond = non_null_cond # 添加当前层级的when条件 condition_list.append(F.when(full_cond, F.lit(col_name))) # 拼接所有when条件 final_condition = reduce(lambda a, b: a.when(b._expr, b._value), condition_list).alias("Hierarchy") new_df = df.select("Dimension1", "Dimention2", final_condition) new_df.show()
两种方法都支持任意数量的层级列,只需修改hierarchies列表即可。
内容的提问来源于stack exchange,提问作者tommyhmt
相关产品推荐
相关产品推荐

