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
相关产品推荐
相关产品推荐

