Spark读取特殊编码转义CSV时collect正常但查询触发解析报错
问题根因
这是Spark CSV数据源的已知解析bug,核心原因是Catalyst优化器的列裁剪下推逻辑,和multiline模式的CSV解析流程存在兼容性缺陷,和业务代码写法、测试CSV文件本身的格式正确性没有关系。
两种操作返回不同结果的逻辑非常清晰:
- 执行
df.collect()时,需要返回所有列的完整数据,不会触发列裁剪优化,CSV解析器会逐行完整跟踪引号闭合、转义字符状态,按照配置的分割规则正确切分三个字段,因此能返回符合预期的解析结果。 - 执行
df.count()、df.select('record_id').distinct().count()这类操作时,Spark会自动做执行计划优化:判断到逻辑只需要用到部分字段,就会把投影规则下推到CSV读取层,要求解析器只提取目标列、跳过其他字段的解析,以此减少解析开销提升性能。问题就出在multiline模式的部分列解析逻辑上:解析器在跳过不需要的字段时,没有正确维护转义引号的闭合状态,错误地把第一行dummy_after字段里逗号后的内容,识别成了record_id列的整数值做类型转换,最终抛出NumberFormatException。
复现必要条件
只有同时满足以下配置和数据特征时才会触发该问题:
- 读取CSV时开启了
multiline=true配置 - 读取时显式指定了StructType表结构,没有开启schema自动推断
- 查询逻辑没有引用全部字段,成功触发了列裁剪下推优化
- CSV内容中存在被双引号包裹、字段值内部包含逗号/转义双引号的记录
临时规避方案
在使用未修复该问题的Spark版本时,可以任选以下方案绕过:
- 读取CSV时追加配置
.option("pushDownProjection", "false"),直接关闭投影下推,强制解析器每次都完整解析整行所有字段,解析行为和collect()完全一致,不会再出现字段错位。 - 读取CSV后先把DataFrame落地为临时Parquet表、或者做一次checkpoint截断执行计划血缘,再基于落地后的中间结果做统计、去重等后续计算,避免直接在原始CSV源上执行带列裁剪的查询。
- 升级到已经合入该bug修复的Spark版本。
注意:这个bug的异常表现存在随机性,如果文件分片策略、数据分布刚好让解析器的状态跟踪没有错位,部分查询可能偶然返回正确结果,但这属于不稳定的偶然行为,生产环境不要依赖这种偶然正确性。
内容的提问来源于stack exchange,提问作者Łukasz Rogalski
相关产品推荐
相关产品推荐

