Spark DataFrame Overwrite写CSV后读取出现groupBy键重复问题
问题原因
- 字符串列存在不可见字符:你观测到的相同
word值实际可能包含前后空格、全角空格、零宽字符、换行符等不可见特殊符号,groupBy算子会判定为不同的分组键,写入CSV后展示时隐藏了特殊字符,就呈现为重复的分组结果。 - 读取时自动类型推导导致值趋同:CSV是纯文本格式,不会保留原始字段类型信息。如果你的
word列取值看起来符合数字特征,读取时Spark默认的自动类型推导会将字符串类型的word转为数字类型,原本字符串层面不同的值(比如带前置空格的" 1001"和无空格的"1001")转成数字后会变成相同值,看起来就是重复的分组键。 - 字段包含特殊字符导致行拆分:如果
word列本身包含你设置的分隔符\t或者换行符,CSV写入时未做转义的情况下,单条数据会被拆分为多行,读取时就会出现非法的重复分组键。
解决方案
你可以按照以下顺序排查解决:
- 先确认分组逻辑本身是否产生重复键:写入
df2前执行以下代码,如果返回结果不为空,说明问题出在分组阶段,和写入读取流程无关:
df2.groupBy($"word").count().filter($"count" > 1).show()
如果确定是分组阶段的问题,对word列做清洗后再分组:
// 去除所有空白字符、前后空格后再分组 val cleanedDf = df1.withColumn("word", trim(regexp_replace($"word", "\\s+", ""))) val df2 = cleanedDf.groupBy($"word").agg(sum($"word_num") as "cnt")
- 如果分组阶段无重复,说明问题出在读写环节,优先推荐更换存储格式为Parquet/ORC,这类结构化存储格式自带schema信息,不会出现类型推导问题,性能也远高于CSV:
// 写入Parquet格式 df2.write.format("parquet") .mode(SaveMode.Overwrite) .save(outputPath) // 读取时自动保留原始schema val pFeatureDF = spark.read.format("parquet").load(outputPath)
- 如果必须使用CSV格式,读取时强制指定schema,避免自动类型推导:
import org.apache.spark.sql.types._ // 定义和写入时完全一致的schema val csvSchema = StructType(Seq( StructField("word", StringType, nullable = true), StructField("cnt", LongType, nullable = true) )) val pFeatureDF = spark.read.format("csv") .option("header", "true") .option("delimiter", "\t") .option("escape", "\"") // 开启转义,避免字段包含分隔符导致行拆分 .schema(csvSchema) .load(outputPath)
内容的提问来源于stack exchange,提问作者withing
相关产品推荐
相关产品推荐

