Spark读取CSV时如何配置支持逗号作为小数分隔符的数字格式
解决方案
Spark 原生 CSV 数据源没有提供单独配置小数分隔符、千分位分组符的独立选项,可根据使用的Spark版本和数据格式选择以下两种方案实现解析,不需要强制只能读取为字符串后手动处理:
方案1:指定locale参数(Spark 2.4及以上版本支持)
从Spark 2.4版本开始,CSV数据源新增locale配置项,用于指定解析时的区域设置,只要设置为使用逗号作为小数分隔符的区域(比如德语de-DE、法语fr-FR),即可直接在读取阶段自动识别逗号格式的小数。
注意:示例代码中给字段定义的类型是
IntegerType,但报错的数值65,9是带小数位的浮点数,就算分隔符格式正确,转整数也会抛出格式错误,必须先将字段类型修正为DoubleType或对应精度的DecimalType。
示例代码:
import org.apache.spark.sql.types._ val df = spark.read.format("csv") .schema(StructType(Seq(StructField("result", DoubleType, true)))) .option("mode", "FAILFAST") .option("delimiter", "|") .option("encoding", "utf8") .option("locale", "de-DE") .load(file)
该方案对符合对应区域标准数值格式的数据(比如带点号千分位的1.234,56格式)也能正常解析,但如果数据的数值格式和对应locale的标准格式不完全匹配,仍然会抛出解析异常。
方案2:读取为字符串后手动转换(全版本兼容,稳定性最高)
如果Spark版本低于2.4,或者数据格式和标准locale格式存在差异,优先使用该方案:读取阶段将数值字段定义为StringType,加载完成后替换逗号为点号,再转换为对应数值类型即可。
如果需要保留FAILFAST模式的严格校验逻辑,可以通过UDF实现异常捕获,遇到非法格式主动抛出错误,和原生解析的行为保持一致。
示例代码:
import org.apache.spark.sql.types._ import org.apache.spark.sql.functions._ // 基础转换实现 val baseDf = spark.read.format("csv") .schema(StructType(Seq(StructField("result", StringType, true)))) .option("mode", "FAILFAST") .option("delimiter", "|") .option("encoding", "utf8") .load(file) .withColumn("result", regexp_replace(col("result"), ",", ".").cast(DoubleType)) // 带严格校验的UDF实现 val parseCommaDecimal = udf((numStr: String) => { if (numStr == null) null else try { numStr.replace(",", ".").toDouble } catch { case e: NumberFormatException => throw new RuntimeException(s"Malformed numeric value: $numStr", e) } }) val validatedDf = baseDf.withColumn("result", parseCommaDecimal(col("result")))
该方案灵活性最高,可以根据实际业务数据的格式(比如同时存在千分位、多余空格等特殊情况)自定义转换逻辑,不会受locale标准格式的限制。
内容的提问来源于stack exchange,提问作者Kombajn zbożowy
相关产品推荐
相关产品推荐

