基于另一Spark DataFrame权重列表填充列(PySpark)
PySpark按权重填充DataFrame列
需求说明
现有两个DataFrame:
- DataFrame A(待填充):
| ID | Name |
|---|---|
| 1 | Null |
| 2 | Null |
| 3 | Null |
| 4 | Null |
- DataFrame B(权重配置):
| Name | weight |
|---|---|
| Max | 0.5 |
| Mike | 0.25 |
| John | 0.25 |
需要按B中Name的权重占比(Max 50%、Mike 25%、John 25%)填充A的Name列,得到目标结果。
实现步骤与代码
1. 初始化并创建示例DataFrame
from pyspark.sql import SparkSession from pyspark.sql.functions import * from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("WeightedFill").getOrCreate() # 创建DataFrame A data_a = [(1, None), (2, None), (3, None), (4, None)] df_a = spark.createDataFrame(data_a, ["ID", "Name"]) # 创建DataFrame B data_b = [("Max", 0.5), ("Mike", 0.25), ("John", 0.25)] df_b = spark.createDataFrame(data_b, ["Name", "weight"])
2. 计算每个Name的填充行数
先获取A的总记录数,结合权重算出每个Name需要填充的行数:
total_rows = df_a.count() # 用ceil确保得到整数行数,避免小数 df_b_with_count = df_b.withColumn("fill_count", ceil(col("weight") * total_rows))
3. 生成权重对应的重复Name序列
将每个Name按计算出的数量重复,生成连续的Name列表:
# 生成重复的Name行 df_repeated = df_b_with_count.withColumn("repeat", explode(array_repeat(col("Name"), col("fill_count")))) name_list = [row["repeat"] for row in df_repeated.select("repeat").collect()]
4. 匹配行号完成填充
给A添加行号,将生成的Name列表与ID对应,替换原Name列:
# 按ID排序添加行号 window = Window.orderBy("ID") df_a_with_row = df_a.withColumn("row_num", row_number().over(window)) # 创建行号与Name的映射表 name_df = spark.createDataFrame(enumerate(name_list, 1), ["row_num", "filled_name"]) # 关联并替换列 final_df = df_a_with_row.join(name_df, on="row_num", how="inner") \ .drop("row_num", "filled_name") \ .withColumn("Name", col("filled_name")) # 展示结果 final_df.show()
运行结果
+---+----+ | ID|Name| +---+----+ | 1| Max| | 2| Max| | 3|Mike| | 4|John| +---+----+
补充说明
- 若A的行数较多,建议避免使用
collect()拉取数据到Driver端,可改用Spark分布式方式生成序列(比如基于权重区间分配行号范围) - 若权重计算后的行数总和与A的总行数有偏差,可调整
ceil为round或floor,或额外处理剩余/缺失的行数
内容的提问来源于stack exchange,提问作者Cri
相关产品推荐
相关产品推荐

