PySpark中将带引号的单列CSV DataFrame转为Delta格式
解决方案
你的问题核心是直接按逗号拆分会破坏双引号内的内容——因为split(',')会把引号里的逗号也当成分隔符,导致列数混乱。针对百万级数据,推荐两种高效且正确的处理方案:
方案一:用from_csv解析(推荐,更稳健高效)
这是PySpark官方推荐的处理带格式CSV内容的方法,完全符合RFC4180标准(自动忽略引号内的分隔符),适合大数据量场景。
步骤1:定义Schema(可选但推荐)
提前定义列的类型能避免Spark自动推断类型的开销,同时保证数据一致性。根据你的示例,Schema可以这样定义:
from pyspark.sql.types import StructType, StructField, IntegerType, StringType schema = StructType([ StructField("col1", IntegerType()), StructField("col2", StringType()), StructField("col3", IntegerType()), StructField("col4", StringType()), # 对应"SAM,K,Clarke" StructField("col5", StringType()), StructField("col6", IntegerType()), StructField("col7", IntegerType()), StructField("col8", StringType()), StructField("col9", IntegerType()), StructField("col10", IntegerType()), StructField("col11", StringType()), # 对应带逗号的邮箱 StructField("col12", IntegerType()), StructField("col13", IntegerType()), StructField("col14", IntegerType()), StructField("col15", StringType()), StructField("col16", IntegerType()), StructField("col17", IntegerType()), StructField("col18", StringType()), # 对应带逗号的地址 StructField("col19", StringType()), StructField("col20", StringType()), StructField("col21", StringType()), StructField("col22", IntegerType()), StructField("col23", StringType()), StructField("col24", IntegerType()), StructField("col25", IntegerType()), # 补充剩余4列的类型(根据实际数据调整) StructField("col26", StringType()), StructField("col27", StringType()), StructField("col28", StringType()), StructField("col29", IntegerType()) ])
步骤2:解析单列内容并展开列
使用from_csv函数按CSV规则解析Details列,再将解析后的结构体展开为独立列:
from pyspark.sql.functions import from_csv, col # 解析CSV内容 df_parsed = df_single.select( from_csv( col("Details"), schema, {"quote": '"', "escape": '"', "sep": ","} # 指定引号和分隔符规则 ).alias("parsed_data") ) # 展开结构体为独立列 df_final = df_parsed.select("parsed_data.*")
方案二:正则表达式拆分(适合快速验证)
用正则匹配不在双引号内的逗号作为分隔符,适合临时验证场景,但大数据量下性能略逊于方案一:
from pyspark.sql.functions import regexp_split, col # 用正则拆分:匹配不在引号内的逗号 df_split = df_single.select( regexp_split( col("Details"), r',(?=(?:[^"]*"[^"]*")*[^"]*$)' # 正则逻辑:仅匹配后面跟着偶数个引号的逗号 ).alias("split_cols") ) # 展开为29列 select_expr = [col("split_cols")[i].alias(f"col{i+1}") for i in range(29)] df_final = df_split.select(*select_expr)
保存为Delta格式
处理完成后,直接写入Delta表:
# 保存到指定路径 df_final.write.format("delta").mode("overwrite").save("/your/delta/path") # 或者保存为Hive表(如果配置了Metastore) # df_final.write.format("delta").mode("overwrite").saveAsTable("your_db.your_table")
为什么你的原方法失效?
split(',')会拆分所有逗号,包括双引号内部的逗号(比如"SAM,K,Clarke"会被拆成3个列),导致实际列数远超过29,后续索引取值自然会出错。
内容的提问来源于stack exchange,提问作者Gourav Joshi
相关产品推荐
相关产品推荐

