PySpark加载分号分隔CSV时&引发列错位的解决方案咨询
解决方案
方案1:RDD预处理后直接转DataFrame(无多余IO,适配150GB大文件场景)
Spark CSV数据源原生支持读取RDD[String]类型的输入,无需将预处理后的RDD落地为中间文件再重新读取,可省掉大量磁盘IO开销,非常适合大文件处理场景:
# 读取原始文件为RDD,直接替换目标字符串 rdd = sc.textFile('/mnt/input/AMP test.csv').map(lambda x: x.replace('&', '&')) # 直接将处理后的RDD加载为CSV格式DataFrame df = spark.read.option("delimiter", ";").option("header","true").csv(rdd) df.show()
该方案逻辑和你当前的临时方案完全一致,但省略了中间文件的写入、读取步骤,是当前场景下性能最优的正规实现方式。
方案2:字段带引号包裹场景可直接用读取参数处理
如果你的原始CSV生成时就用引号包裹了所有字段值,示例格式如下:
ID;FirstName;LastName 1;"Chandler";"Bing" 2;"Ross & Monica";"Geller"
你可以直接指定quote参数让Spark忽略引号内部的分号,避免拆分错位,再配合内置函数处理HTML转义即可:
from pyspark.sql.functions import html_decode df = spark.read \ .option("delimiter", ";") \ .option("header","true") \ .option("quote", '"') \ .csv('/mnt/input/AMP test.csv') # 处理HTML转义,将&转为& df = df.withColumn("FirstName", html_decode("FirstName")) df.show()
如果使用的是Spark 3.0以下版本,没有内置html_decode函数,可以替换为正则替换:
from pyspark.sql.functions import regexp_replace df = df.withColumn("FirstName", regexp_replace("FirstName", "&", "&"))
内容的提问来源于stack exchange,提问作者Connell.O'Donnell
相关产品推荐
相关产品推荐

