【问题标题】:How to make first line of text file as header and skip second line in spark scala如何将文本文件的第一行作为标题并在 spark scala 中跳过第二行
【发布时间】:2020-06-25 11:39:47
【问题描述】:

我试图弄清楚如何使用文本文件的第一行作为标题并跳过秒行。到目前为止,我已经尝试过:

scala> val file = spark.sparkContext.textFile("/home/webwerks/Desktop/UseCase-03-March/Temp/temp.out")  
file: org.apache.spark.rdd.RDD[String] = /home/webwerks/Desktop/UseCase-03-March/Temp/temp.out MapPartitionsRDD[40] at textFile at <console>:23

scala> val clean = file.flatMap(x=>x.split("\t")).filter(x=> !(x.contains("-")))
clean: org.apache.spark.rdd.RDD[String] = MapPartitionsRDD[42] at filter at <console>:25

scala> val df=clean.toDF()
df: org.apache.spark.sql.DataFrame = [value: string]

scala> df.show
+--------------------+
|               value|
+--------------------+
|time         task...|
|03:27:51.199 FCPH...|
|03:27:51.199 PORT...|
|03:27:51.200 PORT...|
|03:27:51.200 PORT...|
|03:27:59.377 PORT...|
|03:27:59.377 PORT...|
|03:27:59.377 FCPH...|
|03:27:59.377 FCPH...|
|03:28:00.468 PORT...|
|03:28:00.468 PORT...|
|03:28:00.469 FCPH...|
|03:28:00.469 FCPH...|
|03:28:01.197 FCPH...|
|03:28:01.197 FCPH...|
|03:28:01.197 PORT...|
|03:28:01.198 PORT...|
|03:28:09.380 PORT...|
|03:28:09.380 PORT...|
|03:28:09.380 FCPH...|

这里我想第一行作为标题,数据应该用制表符分隔

数据是这样的:

time         task       event   port cmd  args
--------------------------------------------------------------------------------------
03:27:51.199 FCPH       seq      13   28  00300000,00000000,00000591,00020182,00000000
03:27:51.199 PORT       Rx       11    0  c0fffffd,00fffffd,0ed10335,00000001
03:27:51.200 PORT       Tx       13   40  02fffffd,00fffffd,0ed3ffff,14000000
03:27:51.200 PORT       Rx       13    0  c0fffffd,00fffffd,0ed329ae,00000001
03:27:59.377 PORT       Rx       15   40  02fffffd,00fffffd,0336ffff,14000000
03:27:59.377 PORT       Tx       15    0  c0fffffd,00fffffd,03360ed2,00000001
03:27:59.377 FCPH       read     15   40  02fffffd,00fffffd,d0000000,00000000,03360ed2
03:27:59.377 FCPH       seq      15   28  22380000,03360ed2,0000052b,0000001c,00000000
03:28:00.468 PORT       Rx       13   40  02fffffd,00fffffd,29afffff,14000000
03:28:00.468 PORT       Tx       13    0  c0fffffd,00fffffd,29af0ed5,00000001

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:
            scala> val ds = spark.read.textFile("data.txt")  > spark-v2.0
                             (or) 
            val ds = spark.sparkContext.textFile("data.txt")
    
            scala> val schemaArr = ds.filter(x=>x.contains("time")).collect.mkString.split("\t").toList
    
            scala> val df = ds.filter(x=> !x.contains("time"))
                              .map(x=>{
                                    val cols = x.split("\t")
                                    (cols(0),cols(1),cols(2),cols(3),cols(4),cols(5))
                                   }).toDF(schemaArr:_*)
    
            scala> df.show(false)
            +------------+----+-----+----+---+--------------------------------------------+
            |time        |task|event|port|cmd|args                                        |
            +------------+----+-----+----+---+--------------------------------------------+
            |03:27:51.199|FCPH|seq  |13  |28 |00300000,00000000,00000591,00020182,00000000|
            |03:27:51.199|PORT|Rx   |11  | 0 |c0fffffd,00fffffd,0ed10335,00000001         |
            |03:27:51.200|PORT|Tx   |13  |40 |02fffffd,00fffffd,0ed3ffff,14000000         |
            |03:27:51.200|PORT|Rx   |13  | 0 |c0fffffd,00fffffd,0ed329ae,00000001         |
            |03:27:59.377|PORT|Rx   |15  |40 |02fffffd,00fffffd,0336ffff,14000000         |
            |03:27:59.377|PORT|Tx   |15  | 0 |c0fffffd,00fffffd,03360ed2,00000001         |
            |03:27:59.377|FCPH|read |15  |40 |02fffffd,00fffffd,d0000000,00000000,03360ed2|
            |03:27:59.377|FCPH|seq  |15  |28 |22380000,03360ed2,0000052b,0000001c,00000000|
            |03:28:00.468|PORT|Rx   |13  |40 |02fffffd,00fffffd,29afffff,14000000         |
            |03:28:00.468|PORT|Tx   |13  | 0 |c0fffffd,00fffffd,29af0ed5,00000001         |
            +------------+----+-----+----+---+--------------------------------------------+
    

    请尝试类似上面的方法,如果你想要模式,然后使用服装模式应用到它

    【讨论】:

    • 它抛出一个错误,即线程“main”java.lang.IllegalArgumentException中的异常:要求失败:列数不匹配。旧列名(6):_1、_2、_3、_4、_5、_6 新列名(1):时间任务事件端口cmd args
    • 你能把你的代码粘贴一次吗?如果它对你有帮助,请告诉我。
    • def main(args:Array[String]) { val spark=SparkSession.builder.master("local").appName("CleaningJson").enableHiveSupport().getOrCreate() 导入火花。 implicits._ val ds = spark.read.textFile("/home/webwerks/Desktop/UseCase-03-March/Temp/temp.out") val schemaArr = ds.filter(x=>x.contains("time" )).collect.mkString.split("\t").toList val df = ds.filter(x=> !x.contains("time")) .map(x=>{val cols = x.split( "\t") (cols(0),cols(1),cols(2),cols(3),cols(4),cols(5)) }).toDF(schemaArr:_*) df.show(假)} }
    • 你能检查一下 schemaArr 的大小,并告诉我,你的代码对我有用。 scala> schemaArr.size res0: Int = 6 other wisepilase 提供正确的输入数据
    • scala> schemaArr.size res6: Int = 1 实际上它的分隔符不是 \t 我猜那我怎么能用 \t 替换任何空格然后转换成数据框
    猜你喜欢
    • 2015-01-09
    • 2020-12-15
    • 1970-01-01
    • 2013-03-29
    • 2021-11-05
    • 2020-06-29
    • 2017-03-04
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多