【问题标题】:Spark read key value pairs from the file into a DataframeSpark将文件中的键值对读入Dataframe
【发布时间】:2020-11-09 06:22:33
【问题描述】:

我需要读取一个日志文件并将其转换为 spark 数据帧。

输入文件内容:

dateCreated   : 20200720
customerId    :  001
dateCreated   : 20200720
customerId    :  002
dateCreated   : 20200721
customerId    :  003

预期的数据框:

---------------------------
|dateCreated | customerId |
---------------------------
|20200720    | 001        |
|20200720    | 002        |
|20200721    | 003        |
|------------|------------|

火花代码:

val spark = org.apache.spark.sql.SparkSession.builder.master("local").getOrCreate
    val inputFile = "C:\\log_data.txt"
    val rddFromFile = spark.sparkContext.textFile(inputFile)

    val rdd = rddFromFile.map(f => {
      f.split(":")
    })

    rdd.foreach(f => {
      println(f(0) + "\t" + f(1))
    })

关于如何将上述 rdd 转换为所需的 DF 的任何想法?

【问题讨论】:

    标签: python scala apache-spark pyspark apache-spark-sql


    【解决方案1】:

    检查下面的代码。

    scala> "cat /tmp/sample/input.csv".!
    dateCreated   : 20200720
    customerId    :  001
    dateCreated   : 20200720
    customerId    :  002
    dateCreated   : 20200721
    customerId    :  003
    
    scala> val df = spark.read.text("/tmp/sample").select(split($"value",":").as("data"))
    df: org.apache.spark.sql.DataFrame = [data: array<string>]
    
    scala> df.show(false)
    +---------------------------+
    |data                       |
    +---------------------------+
    |[dateCreated   ,  20200720]|
    |[customerId    ,   001]    |
    |[dateCreated   ,  20200720]|
    |[customerId    ,   002]    |
    |[dateCreated   ,  20200721]|
    |[customerId    ,   003]    |
    +---------------------------+
    
    scala> import org.apache.spark.sql.expressions._
    import org.apache.spark.sql.expressions._
    
    scala> val windowSpec = Window.orderBy($"id".asc)
    
    scala> df
    .select(trim($"data"(0)).as("data"),trim($"data"(1)).as("values"))
    .select(map($"data",$"values").as("data"))
    .select($"data"("dateCreated").as("dateCreated"),$"data"("customerId").as("customerId"))
    .withColumn("id",monotonically_increasing_id)
    .withColumn("customerId",lead($"customerId",1).over(windowSpec))
    .where($"customerId".isNotNull)
    .drop("id")
    .show(false)
    
    +-----------+----------+
    |dateCreated|customerId|
    +-----------+----------+
    |20200720   |001       |
    |20200720   |002       |
    |20200721   |003       |
    +-----------+----------+
    

    【讨论】:

    • 除了使用窗口函数还有其他方法吗? .我希望窗口函数是昂贵的操作并且输入文件很大。
    猜你喜欢
    • 2022-01-13
    • 1970-01-01
    • 2017-04-25
    • 2012-09-20
    • 2020-02-02
    • 2016-11-29
    • 1970-01-01
    • 2021-05-05
    • 1970-01-01
    相关资源
    最近更新 更多