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

Spark中按|分割字符串及RDD转DataFrame时数据损坏问题

问题根源与解决方案

你遇到的问题核心在于**String.split()方法接受的是正则表达式参数**,而竖线|在正则表达式里是一个特殊元字符(表示“或”的逻辑),直接用split("|")会导致正则引擎把每个字符都当作分隔符,所以才会把字符串拆成单个字符的数组。

修正步骤

1. 修复RDD的分割逻辑

把分割符改成转义后的形式,在Scala字符串中,需要用双反斜杠\\来转义正则里的特殊字符:

import org.apache.spark.sql.Encoder
import spark.implicits._
case class Product(productId:Int, price:Double, saleEvent:String, rivalName:String, fetchTS:String)

val rdd = spark.sparkContext.textFile("/home/prabhat/Documents/Spark/sampledata/competitor_data_10.txt")
// 移除表头
val x = rdd.mapPartitionsWithIndex{(idx,iter) => if(idx==0)iter.drop(1) else iter}
// 用转义后的竖线分割
val splitRdd = x.map(_.split("\\|"))
splitRdd.take(10)

或者更稳妥的方式,用Pattern.quote()来包裹分隔符,自动处理所有正则特殊字符:

import java.util.regex.Pattern
val splitRdd = x.map(_.split(Pattern.quote("|")))

2. 转换为正确的DataFrame

分割后需要把RDD映射到Product样例类,再转成DataFrame:

val productDF = splitRdd.map(arr => Product(
    arr(0).trim.toInt,
    arr(1).trim.toDouble,
    arr(2).trim,
    arr(3).trim,
    arr(4).trim
)).toDF()

productDF.show()

更简便的方案:直接用Spark SQL读取

其实不需要手动处理RDD,Spark SQL的read API可以直接指定分隔符读取这类文件,更简洁且不易出错:

val productDF = spark.read
    .option("sep", "|")
    .option("header", "true")
    .option("inferSchema", "true")
    .csv("/home/prabhat/Documents/Spark/sampledata/competitor_data_10.txt")

productDF.show()

这个方法会自动处理表头、分隔符,还能自动推断字段类型,省去手动转换的麻烦。

验证结果

修正后你会得到预期的输出:

+---------+-------+---------+---------------+-------------------+
|productId|  price|saleEvent|      rivalName|            fetchTS|
+---------+-------+---------+---------------+-------------------+
|    12345|  78.73|  Special|   VistaCart.com|2017-05-11 15:39:30|
|    12345|  45.52|  Regular|ShopYourWay.com|2017-05-11 16:09:43|
|    12345|  89.52|     Sale|MarketPlace.com|2017-05-11 16:07:29|
|      678|1348.73|  Regular|   VistaCart.com|2017-05-11 15:58:06|
+---------+-------+---------+---------------+-------------------+

内容的提问来源于stack exchange,提问作者tryingSpark

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:41:08