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

PySpark中实现基于列的迭代Join直至新增列全为Null

迭代关联PySpark DataFrame实现层级展开

问题背景

给定如下PySpark DataFrame,需要迭代将每一列的末尾值与start列匹配,生成下一级的endN列,直到新增列全为Null为止:

data = [('s1', 's2'),
       ('s1', 's3'),
       ('s2', 's4'),
       ('s3', 's5'),
       ('s5', 's6')]
sdf = spark.createDataFrame(data, schema=['start', 'end'])

解决方案

通过循环执行左关联,每次用最新生成的列匹配原始数据的start字段,生成下一级列,直到新列无有效数据时停止:

# 提取原始数据的映射关系,用于后续关联
lookup_df = sdf.selectExpr("start as lookup_start", "end as lookup_end")

current_df = sdf
current_col = "end"
col_index = 2

while True:
    new_col = f"end{col_index}"
    # 左关联获取下一级的end值
    current_df = current_df.join(
        lookup_df,
        current_df[current_col] == lookup_df["lookup_start"],
        how="left"
    ).drop("lookup_start").withColumnRenamed("lookup_end", new_col)
    
    # 检查新列是否存在非空值
    has_valid_data = current_df.where(f"{new_col} IS NOT NULL").count() > 0
    if not has_valid_data:
        # 全为空则删除该列并终止循环
        current_df = current_df.drop(new_col)
        break
    
    # 更新当前关联列和索引
    current_col = new_col
    col_index += 1

# 查看最终结果
current_df.show()

执行结果

+-----+---+----+----+
|start|end|end2|end3|
+-----+---+----+----+
|   s1| s2|  s4|null|
|   s1| s3|  s5|  s6|
|   s2| s4|null|null|
|   s3| s5|  s6|null|
|   s5| s6|null|null|
+-----+---+----+----+

内容的提问来源于stack exchange,提问作者Rory

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 19:55:18