【问题标题】:Using groupBy in Spark and getting back to a DataFrame在 Spark 中使用 groupBy 并返回 DataFrame
【发布时间】:2016-02-13 16:53:22
【问题描述】:

在使用 Scala 处理 Spark 中的数据帧时,我遇到了困难。如果我有一个要提取一列唯一条目的数据框,那么当我使用 groupBy 时,我不会得到一个数据框。

例如,我有一个名为 logs 的 DataFrame,其格式如下:

machine_id  | event     | other_stuff
 34131231   | thing     |   stuff
 83423984   | notathing | notstuff
 34131231   | thing    | morestuff

并且我希望将事件存储在新的DataFrame 中的唯一机器 ID 允许我进行某种过滤。使用

val machineId = logs
  .where($"event" === "thing")
  .select("machine_id")
  .groupBy("machine_id")

我得到一个分组数据的 val,使用起来很麻烦(或者我不知道如何正确使用这种对象)。得到这个唯一机器 ID 的列表后,我想用它来过滤另一个 DataFrame 以提取单个机器 ID 的所有事件。

我可以看到我想要相当定期地做这种事情,基本的工作流程是:

  1. 从日志表中提取唯一 ID。
  2. 使用唯一 ID 提取特定 ID 的所有事件。
  3. 对已提取的数据进行某种分析。

这是前两个步骤,我希望得到一些指导。

我很欣赏这个例子有点做作,但希望它能解释我的问题。可能是我对GroupedData 对象不够了解,或者(正如我希望的那样)我在数据框中遗漏了一些使这变得容易的东西。我正在使用基于 Scala 2.10.4 构建的 spark 1.5。

谢谢

【问题讨论】:

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


    【解决方案1】:

    只使用distinct 而不是groupBy

    val machineId = logs.where($"event"==="thing").select("machine_id").distinct
    

    相当于 SQL:

    SELECT DISTINCT machine_id FROM logs WHERE event = 'thing'
    

    GroupedData 不能直接使用。它提供了多种方法,其中agg是最通用的,可用于应用不同的聚合函数并将其转换回DataFrame。就 SQL 而言,wheregroupBy 之后的内容相当于这样

    SELECT machine_id, ... FROM logs WHERE event = 'thing' GROUP BY machine_id
    

    其中... 必须由agg 或等效方法提供。

    【讨论】:

      【解决方案2】:

      spark 中的 group by 然后聚合,然后 select 语句将返回一个数据帧。对于您的示例,它应该类似于:

      val machineId = logs
          .groupBy("machine_id", "event")
          .agg(max("other_stuff") )
          .select($"machine_id").where($"event" === "thing")
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2016-04-03
        • 1970-01-01
        • 2018-01-31
        • 2016-05-31
        相关资源
        最近更新 更多