Spark 2.3.0中如何加载含多行记录的CSV文件?
解决Spark 2.3.0读取CSV行错位的问题
嘿,我之前在Spark 2.3.0里处理CSV的时候也碰到过一模一样的行错位问题,尤其是开了multiLine=True之后反而更乱。咱们从几个常见的坑入手,一步步解决:
先排查multiLine的使用场景
multiLine=True是用来处理字段值内包含换行符的CSV(比如Body字段里的换行),但前提是这些带换行的字段必须被引号包裹。如果你的CSV里有未被引号包裹的换行,Spark会把它当成新行的分隔符,直接导致行错位。建议你先打开CSV文件确认:所有带换行的字段是不是都被双引号(或其他引号)包起来了?
别依赖inferSchema,手动指定更可靠
Spark 2.3.0的inferSchema在multiLine模式下很容易“懵”——比如字段里有特殊字符、换行或者空值时,它会错误推断字段类型,进而引发解析错位。你可以先把inferSchema=False,用字符串类型读取所有字段,看看行是不是对齐了:
answer_df = sparkSession.read.csv( './stacksample/Answers_sample.csv', header=True, inferSchema=False, multiLine=True ) answer_df.show(2)
如果这时候行对齐了,那问题肯定出在schema推断上,接下来就手动定义准确的schema:
from pyspark.sql.types import StructType, StructField, IntegerType, StringType, TimestampType # 根据你的CSV实际字段结构定义schema(这里假设是Stack Overflow Answers的结构) answer_schema = StructType([ StructField("Id", IntegerType(), nullable=True), StructField("OwnerUserId", IntegerType(), nullable=True), StructField("CreationDate", TimestampType(), nullable=True), StructField("ParentId", IntegerType(), nullable=True), StructField("Score", IntegerType(), nullable=True), StructField("Body", StringType(), nullable=True) ]) # 用指定的schema读取CSV answer_df = sparkSession.read.csv( './stacksample/Answers_sample.csv', header=True, schema=answer_schema, multiLine=True, quote='"', # 明确指定包裹字段的引号字符,默认是双引号 escape='"' # 处理字段内的双引号(比如CSV里用""表示一个") ) answer_df.show(2)
其他可能的排查点
- 分隔符问题:如果你的CSV不是用逗号分隔(比如制表符),一定要加上
sep='\t'参数,否则Spark会把整个行当成一个字段。 - 文件编码问题:如果CSV是UTF-8带BOM的格式,Spark 2.3.0可能解析错误,试试指定
encoding='UTF-8-SIG'。 - Header行异常:如果CSV的Header行本身包含换行,
header=True会把Header解析成多行,导致后续行错位。这种情况可以先把header=False,然后手动给DataFrame指定列名:answer_df = sparkSession.read.csv( './stacksample/Answers_sample.csv', header=False, schema=answer_schema, multiLine=True ) # 如果第一行是Header,跳过它 answer_df = answer_df.filter(answer_df.Id != "Id")
按照这个流程排查下来,基本能解决行错位的问题啦!
内容的提问来源于stack exchange,提问作者Gaurav Gupta
相关产品推荐
相关产品推荐

