使用PySpark读取CSV文件时疑似编码问题导致数据读取错误
问题复现
- 为完成大学课程,使用pyspark-notebook docker镜像,启动命令如下:
docker pull jupyter/pyspark-notebook docker run -it --rm -p 8888:8888 -v /path/to/my/working/directory:/home/jovyan/work jupyter/pyspark-notebook
- 初始读取CSV代码:
import pyspark from pyspark.sql import SparkSession from pyspark.sql.types import * sc = pyspark.SparkContext('local[*]') spark = SparkSession(sc) spark listings_df = spark.read.csv("listings.csv", header=True, mode='DROPMALFORMED') # 在上一行代码添加encoding="utf8"参数也无法解决问题 listings_df.printSchema()
- 异常表现:
- Spark读取得到16494行,
pandas.read_csv()校验正确行数为16478行 room_type字段出现大量非法值,合法取值仅为['Private room', 'Entire home/apt', 'Hotel room', 'Shared room'],分组统计结果如下:
- Spark读取得到16494行,
+---------------+-----+ | room_type|count| +---------------+-----+ | 169| 1| | 4.88612| 1| | 4.90075| 1| | Shared room| 44| | 35| 1| | 187| 1| | null| 16| | 70| 1| | 27| 1| | 75| 1| | Hotel room| 109| | 198| 1| | 60| 1| | 280| 1| |Entire home/apt|12818| | 220| 1| | 190| 1| | 156| 1| | 450| 1| | 4.88865| 1| +---------------+-----+ only showing top 20 rows
- 已知信息:
- Spark版本为v3.1.2,运行模式
local[*] - 文件编码为UTF-8
- 数据集为Airbnb公开房源统计数据
- Spark版本为v3.1.2,运行模式
解决方案
问题核心原因是数据集内部分文本字段被双引号包裹且包含换行符,Spark默认CSV解析器不会识别双引号内的换行属于同一行,导致单条数据被拆分为多行读取,出现字段错位、非法值、行数异常增多的问题。
直接修改CSV读取配置,新增multiLine参数即可解决,修改后代码如下:
listings_df = spark.read.csv( "listings.csv", header=True, mode='DROPMALFORMED', multiLine=True, # 识别双引号包裹的跨行内容 escape='"' # 明确指定双引号为转义字符 )
修改后可验证:读取行数与pandas读取结果一致,room_type字段仅出现4类合法取值,无异常数值。
内容的提问来源于stack exchange,提问作者maria
相关产品推荐
相关产品推荐

