【问题标题】:RDD/Scala Get one column from RDDRDD/Scala 从 RDD 中获取一列
【发布时间】:2017-02-27 22:01:20
【问题描述】:

我有一个包含多个字段 (username,content,date,bytes) 的 RDD[Log] 文件,我想为每个字段/列找到不同的内容。

例如,我想获取在 RDD 中找到的最小/最大和平均字节数。当我这样做时:

val q1 = cleanRdd.filter(x => x.bytes != 0)

我得到了带有字节的 RDD 的完整行!= 0。但是我怎样才能真正对它们求和、计算平均值、找到最小值/最大值等?如何仅从我的 RDD 中获取一列并对其应用转换?

编辑:Prasad 告诉我有关将类型更改为 dataframe 的事情,但他没有提供有关如何执行此操作的说明,而且我在网站上找不到可靠的答案。任何帮助都会很棒。

编辑:日志类:

case class Log (username: String, date: String, status: Int, content: Int)

使用 cleanRdd.take(5).foreach(println) 会得到类似的结果

Log(199.72.81.55 ,01/Jul/1995:00:00:01 -0400,200,6245)
Log(unicomp6.unicomp.net ,01/Jul/1995:00:00:06 -0400,200,3985)
Log(199.120.110.21 ,01/Jul/1995:00:00:09 -0400,200,4085)
Log(burger.letters.com ,01/Jul/1995:00:00:11 -0400,304,0)
Log(199.120.110.21 ,01/Jul/1995:00:00:11 -0400,200,4179)

【问题讨论】:

    标签: scala apache-spark rdd


    【解决方案1】:

    嗯...你有很多问题。

    所以...您有以下 Log 抽象

    case class Log (username: String, date: String, status: Int, content: Int, byte: Int)
    

    Que - 我怎样才能从我的 RDD 中只取一列。

    Ans - 你有一个带有 RDD 的 map 函数。所以对于RDD[A]map 采用A => B 类型的映射/转换函数将其转换为RDD[B]

    val logRdd: RDD[Log] = ...
    
    val byteRdd = logRdd
      .filter(l => l.bytes != 0)
      .map(l => l.byte)
    

    阙 - 我怎样才能真正总结它们?

    Ans - 您可以使用 reduce / fold / aggregate 来做到这一点。

    val sum = byteRdd.reduce((acc, b) => acc + b)
    
    val sum = byteRdd.fold(0)((acc, b) => acc + b)
    
    val sum = byteRdd.aggregate(0)(
      (acc, b) => acc + b,
      (acc1, acc2) => acc1 + acc2
    )
    

    注意 :: 这里要注意的重要一点是,Int 的总和可以增长到比 Int 可以处理的更大。所以在大多数现实生活中,我们至少应该使用Long 作为累加器,而不是Int,它实际上删除了reducefold 作为选项。我们将只剩下一个聚合。

    val sum = byteRdd.aggregate(0l)(
      (acc, b) => acc + b,
      (acc1, acc2) => acc1 + acc2
    )
    

    现在,如果您必须计算诸如 min、max、avg 之类的多个值,那么我建议您在单个 aggregate 中计算它们,而不是像这样的多个,

    // (count, sum, min, max)
    val accInit = (0, 0, Int.MaxValue, Int.MinValue)
    
    val (count, sum, min, max) = byteRdd.aggregate(accInit)(
      { case ((count, sum, min, max), b) => 
          (count + 1, sum + b, Math.min(min, b), Math.max(max, b)) },
      { case ((count1, sum1, min1, max1), (count2, sum2, min2, max2)) => 
          (count1 + count2, sum1 + sum2, Math.min(min1, min2), Math.max(max1, max2)) }
    })
    
    val avg = sum.toDouble / count
    

    【讨论】:

      【解决方案2】:

      查看DataFrame API。您需要将 RDD 转换为 DataFrame,然后您可以使用 min、max、avg 函数,如下所示:

      val rdd = cleanRdd.filter(x => x.bytes != 0)
      val df = sparkSession.sqlContext.createDataFrame(rdd, classOf[Log])
      

      假设您想对列 bytes 进行操作

      import org.apache.spark.sql.functions._
      
      df.select(avg("bytes")).show
      df.select(min("bytes")).show
      df.select(max("bytes")).show
      

      更新:

      在 spark-shell 中尝试了以下内容。检查屏幕截图以了解结果...

      case class Log (username: String, date: String, status: Int, content: Int)
      
      val inputRDD = sc.parallelize(Seq(Log("199.72.81.55","01/Jul/1995:00:00:01 -0400",200,6245), Log("unicomp6.unicomp.net","01/Jul/1995:00:00:06 -0400",200,3985), Log("199.120.110.21","01/Jul/1995:00:00:09 -0400",200,4085), Log("burger.letters.com","01/Jul/1995:00:00:11 -0400",304,0), Log("199.120.110.21","01/Jul/1995:00:00:11 -0400",200,4179)))
      
      val rdd = inputRDD.filter(x => x.content != 0)
      
      val df = rdd.toDF("username", "date", "status", "content")
      
      df.printSchema
      
      import org.apache.spark.sql.functions._
      
      df.select(avg("content")).show
      df.select(min("content")).show
      df.select(max("content")).show
      

      【讨论】:

      • 有没有办法在不转换我的 RDD 的情况下做到这一点?通常将 RDD 转换为 Dataframe 以执行此类操作吗?
      • 另外,我如何使用你给我的命令,我得到一个错误“错误:未找到:值 sparkSession”
      • 使用SparkSession 实例将RDD 转换为DataFrame。如果您在spark-shell 中运行它,那么请使用sparkSession 而不是spark
      • WARN ObjectStore:在 Metastore 中找不到版本信息。 hive.metastore.schema.verification 未启用,因此记录架构版本 1.2.0 17/02/27 12:21:07 WARN ObjectStore: 无法获取数据库默认值,返回 NoSuchObjectException df: org.apache.spark.sql.DataFrame = [] 这就是我得到的
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-05-06
      • 1970-01-01
      • 1970-01-01
      • 2021-10-16
      • 2017-02-17
      • 2017-10-07
      • 1970-01-01
      相关资源
      最近更新 更多