You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.20 16:04:54