You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark DataFrame Overwrite写CSV后读取出现groupBy键重复问题

问题原因
  • 字符串列存在不可见字符:你观测到的相同word值实际可能包含前后空格、全角空格、零宽字符、换行符等不可见特殊符号,groupBy算子会判定为不同的分组键,写入CSV后展示时隐藏了特殊字符,就呈现为重复的分组结果。
  • 读取时自动类型推导导致值趋同:CSV是纯文本格式,不会保留原始字段类型信息。如果你的word列取值看起来符合数字特征,读取时Spark默认的自动类型推导会将字符串类型的word转为数字类型,原本字符串层面不同的值(比如带前置空格的" 1001"和无空格的"1001")转成数字后会变成相同值,看起来就是重复的分组键。
  • 字段包含特殊字符导致行拆分:如果word列本身包含你设置的分隔符\t或者换行符,CSV写入时未做转义的情况下,单条数据会被拆分为多行,读取时就会出现非法的重复分组键。
解决方案

你可以按照以下顺序排查解决:

  1. 先确认分组逻辑本身是否产生重复键:写入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")
  1. 如果分组阶段无重复,说明问题出在读写环节,优先推荐更换存储格式为Parquet/ORC,这类结构化存储格式自带schema信息,不会出现类型推导问题,性能也远高于CSV:
// 写入Parquet格式
df2.write.format("parquet")
  .mode(SaveMode.Overwrite)
  .save(outputPath)
// 读取时自动保留原始schema
val pFeatureDF = spark.read.format("parquet").load(outputPath)
  1. 如果必须使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.23 23:54:01