如何使用PySpark将报表工具导出的交叉表CSV转换为标准表格
PySpark实现ID空值向前填充方案
核心思路
利用Spark窗口函数的last()方法,跳过空值取最近的非空ID值填充后续空行,由于Spark分布式框架没有默认行顺序,需要先生成单调递增的行号列保证原始CSV的行顺序不变。
完整实现代码
# 导入依赖模块 from pyspark.sql import SparkSession from pyspark.sql.functions import last, monotonically_increasing_id, when, col from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("fill_group_null_id").getOrCreate() # 读取原始CSV文件 df = spark.read.csv( "path/to/your/input.csv", header=True, # 第一行为表头 inferSchema=True # 自动推导字段类型 ) # 可选:如果CSV中空ID是空白字符串而非null,先转换为null df = df.withColumn("ID", when(col("ID") == "", None).otherwise(col("ID"))) # 生成单调递增行号,保证原始行顺序 df_with_row = df.withColumn("row_id", monotonically_increasing_id()) # 定义窗口:按行号升序排序,范围从数据集开头到当前行 window_spec = Window.orderBy("row_id").rowsBetween(Window.unboundedPreceding, 0) # 填充ID列空值:取最近的非空ID值 df_filled = df_with_row.withColumn("ID", last("ID", ignorenulls=True).over(window_spec)) # 删除临时行号列,得到最终结果 df_result = df_filled.drop("row_id") # 验证结果 df_result.show() # 导出为标准CSV df_result.write.csv( "path/to/your/output.csv", header=True, mode="overwrite" )
关键说明
monotonically_increasing_id()生成的行号是全局唯一且单调递增的,完全匹配原始CSV的行顺序,避免分布式读取导致的行乱序问题last()函数的ignorenulls=True参数会跳过所有空值,始终取最近的非空ID值,刚好匹配行分组导出的CSV规则- 如果你的原始ID列是数值类型,填充后可以自行用
cast()方法转换回对应类型
内容的提问来源于stack exchange,提问作者wilson_smyth
相关产品推荐
相关产品推荐

