Spark Scala:使用zipWithIndex移除前2行报错,求解决方案
解决Spark Scala中zipWithIndex过滤表头的错误问题
你的代码报错是因为Scala里访问元组元素的语法用错了——zipWithIndex返回的是(String, Long)类型的元组,不能用row[0]这种方式访问元素,得用._1、._2来分别取元组的第一个和第二个值。
正确的RDD写法
val data = spark.sparkContext.textFile(filename) // 过滤索引大于1的行(移除前两行表头),然后提取行内容 val filteredData = data.zipWithIndex() .filter(row => row._2 > 1) .map(row => row._1)
代码说明
row._2就是zipWithIndex分配的索引(从0开始),过滤条件row._2 > 1会保留索引为2及以后的行,正好去掉前两行表头。- 最后用
map(row => row._1)把元组里的行内容提取出来,得到的就是去掉表头后的纯数据。
如果用DataFrame API的话,写法会更直观,也可以参考:
import spark.implicits._ val df = spark.read.text(filename) val filteredDf = df.withColumn("index", monotonically_increasing_id()) .filter($"index" > 1) .drop("index")
验证结果
用你的示例数据测试,运行后会得到:
1,123,ramakrishna 2,234,deepthi
内容的提问来源于stack exchange,提问作者Anonymous
相关产品推荐
相关产品推荐

