如何将PySpark读取CSV得到的单列累积数据拆分为多列DataFrame
可行的简便实现方案
方案1:调整CSV读取参数(最优)
你当前读取CSV时未配置分隔符、表头参数,才会出现所有数据合并到_c0单列的情况,直接调整读取配置即可一步得到目标DataFrame,无需后续加工:
df = spark.read.format("csv") .option("sep", " ") # 指定分隔符为空格 .option("header", "true") # 读取第一行作为列名 .load(path) .select("A", "B", "C", "D") # 仅保留需要的列,自动过滤多余的E列
方案2:split函数拆分单列(适用于已读取到单列DataFrame的场景)
如果已经拿到只有_c0单列的df1,用split函数按空格拆分行内容即可,不需要记忆每个字段的起始位置和长度:
from pyspark.sql.functions import split, trim, col # 提取表头行内容 header_row = df1.first()["_c0"] # 过滤表头行后拆分内容,提取对应列 outcs = df1.filter(col("_c0") != header_row) .withColumn("cols", split("_c0", " ")) .select( trim(col("cols")[0]).alias("A"), trim(col("cols")[1]).alias("B"), trim(col("cols")[2]).alias("C"), trim(col("cols")[3]).alias("D") )
方案3:固定宽度格式专用读取方案
如果源文件是固定宽度格式(每个字段占固定字符数),可以用spark固定宽度数据源直接读取,无需手动编写多个substring逻辑:
df = spark.read.format("fixed-width") .option("header", "true") .option("columnWidths", [10, 36, 36, 36]) # 对应你用到的各字段长度,自动计算偏移量 .load(path) .select("A", "B", "C", "D")
内容的提问来源于stack exchange,提问作者Shankhadip Kundu
相关产品推荐
相关产品推荐

