Spark读取含特殊字符的CSV文件:列偏移问题求解
解决Spark读取含特殊字符CSV的列偏移问题
问题背景
要读取gzip压缩格式的CSV文件,分隔符为逗号(,),其中Description列包含逗号、反斜杠等特殊字符。现有Spark读取配置存在矛盾:
- 仅用默认CSV配置时,第1行(含逗号的
Description)能正常解析,但第2行的Description显示为逗号且列偏移; - 添加
escape配置后,第2行恢复正常,但第1行因Description里的逗号被当成分隔符,导致列偏移、Amount列数据错乱。
且无法修改源文件或分隔符。
解决方案:自定义行解析逻辑
由于源文件未遵循标准CSV规范(含分隔符的字段未用引号包裹),Spark默认CSV reader无法同时兼容两种特殊场景,需通过逐行自定义解析来解决:
步骤1:定义DataFrame Schema
先明确目标Schema,避免自动推断出错:
from pyspark.sql.types import StructType, StructField, StringType # 匹配源文件的列结构 schema = StructType([ StructField("Id", StringType(), nullable=True), StructField("Description", StringType(), nullable=True), StructField("No", StringType(), nullable=True), StructField("Amount", StringType(), nullable=True) ])
步骤2:以文本模式读取文件并自定义解析
跳过表头后,针对每一行按列的固定位置拆分:
- 第一列是
Id,取第一个逗号前的内容; - 最后两列是
No和Amount,从末尾倒推两个逗号拆分; - 中间所有内容归为
Description,确保其内部的逗号不会影响列拆分。
代码示例:
# 以文本模式读取gzip文件,text模式支持compression选项 raw_rdd = spark.sparkContext.textFile(path_dm) # 分离表头和数据行 header_line = raw_rdd.first() data_lines_rdd = raw_rdd.filter(lambda line: line != header_line) def parse_single_line(line): # 拆分第一列Id first_comma_pos = line.find(',') if first_comma_pos == -1: return (line, None, None, None) id_val = line[:first_comma_pos] remaining_content = line[first_comma_pos + 1:] # 拆分最后一列Amount last_comma_pos = remaining_content.rfind(',') if last_comma_pos == -1: return (id_val, remaining_content, None, None) amount_val = remaining_content[last_comma_pos + 1:] content_before_amount = remaining_content[:last_comma_pos] # 拆分倒数第二列No,中间部分为Description second_last_comma_pos = content_before_amount.rfind(',') if second_last_comma_pos == -1: return (id_val, content_before_amount, None, amount_val) no_val = content_before_amount[second_last_comma_pos + 1:] description_val = content_before_amount[:second_last_comma_pos] # 还原反斜杠字符 if description_val: description_val = description_val.replace('\\', '\\') return (id_val, description_val, no_val, amount_val) # 解析后转成DataFrame parsed_rdd = data_lines_rdd.map(parse_single_line) inputFileDf1 = spark.createDataFrame(parsed_rdd, schema=schema)
方案说明
这个逻辑利用了列的固定位置特性,完全规避了默认CSV reader对分隔符和转义符的冲突处理:
- 不管
Description里有多少逗号,都会被完整保留; - 反斜杠字符通过字符串替换还原,不会被当成转义符干扰列拆分;
- 最后两列
No和Amount始终能准确定位,避免偏移。
内容的提问来源于stack exchange,提问作者Priyam
相关产品推荐
相关产品推荐

