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
相关产品推荐
相关产品推荐

