如何使用PySpark实现分组内非空值填充同组空字符串
PySpark分组内空值填充实现方案
实现思路
利用Spark窗口函数按group列分组,提取每个分组内的唯一非空值,填充该分组下所有空行即可,具体步骤如下:
- 先将
values_to_copy列中的空字符串转换为null,适配Spark内置函数的空值处理逻辑 - 按
group字段创建窗口分区 - 提取每个分组内的非空值,覆盖所有行的空值
可运行代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import col, first, when from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("group_fill").getOrCreate() # 示例数据构造 data = [ ("value_1", 31, 1), ("", 31, 1), ("value_1", 31, 1), ("value_2", 152, 1), ("", 152, 1), ("", 153, 1), ("value_3", 153, 1), ("value_4", 154, 1), ("", 154, 1), ] df = spark.createDataFrame(data, schema=["values_to_copy", "group", "flag"]) # 1. 空字符串转null df = df.withColumn("values_to_copy", when(col("values_to_copy") == "", None).otherwise(col("values_to_copy"))) # 2. 定义按group分组的窗口 window_spec = Window.partitionBy("group") # 3. 分组内取第一个非空值填充所有行 df = df.withColumn("values_to_copy", first("values_to_copy", ignorenulls=True).over(window_spec)) # 输出结果 df.show()
补充说明
如果你的业务场景中单个分组存在多个不同的非空values_to_copy值,可以在窗口定义中新增排序规则,结合last/max等函数匹配你的填充逻辑。
内容的提问来源于stack exchange,提问作者Juan Rodríguez
相关产品推荐
相关产品推荐

