【问题标题】:Counting both occurrences and cooccurrences in a DF计算 DF 中的出现和同时​​出现
【发布时间】:2016-07-29 06:40:51
【问题描述】:

我想计算我在 Spark 数据框中的两个变量 xy 之间的 mutual information (MI),如下所示:

scala> df.show()
+---+---+
|  x|  y|
+---+---+
|  0| DO|
|  1| FR|
|  0| MK|
|  0| FR|
|  0| RU|
|  0| TN|
|  0| TN|
|  0| KW|
|  1| RU|
|  0| JP|
|  0| US|
|  0| CL|
|  0| ES|
|  0| KR|
|  0| US|
|  0| IT|
|  0| SE|
|  0| MX|
|  0| CN|
|  1| EE|
+---+---+

在我的例子中,x 恰好是事件是否正在发生 (x = 1) 或不发生 (x = 0),y 是国家代码,但这些变量可以代表任何东西。要计算 xy 之间的 MI,我希望将上述数据框按 x, y 对分组,并添加以下三列:

  • x 的出现次数
  • y 的出现次数
  • x, y 的出现次数

在上面的简短示例中,它看起来像

x, y, count_x, count_y, count_xy
0, FR, 17, 2, 1
1, FR, 3, 2, 1
...

然后我只需要计算每个 x, y 对的互信息项并将它们求和。

到目前为止,我已经能够按 x, y 对分组并聚合 count(*) 列,但我找不到添加 xy 计数的有效方法。我目前的解决方案是将 DF 转换为数组并手动计算出现次数和同时出现次数。当y 是一个国家时,它运行良好,但当y 的基数变大时,它需要很长时间。关于如何以更 Sparkish 的方式做到这一点的任何建议?

提前致谢!

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    我会使用 RDD,为每个用例生成一个键,按键计数并加入结果。这样我就知道具体是什么阶段了。

    rdd.cache() // rdd is your data [x,y]
    val xCnt:RDD[Int, Int] = rdd.countByKey
    val yCnt:RDD[String, Int] = rdd.countByValue
    val xyCnt:RDD[(Int,String), Int] = rdd.map((x, y) => ((x,y), x,y)).countByKey
    val tmp = xCnt.cartsian(yCnt).map(((x, xCnt),(y, yCnt)) => ((x,y),xCnt,yCnt))
    val miReady = tmp.join(xyCnt).map(((x,y), ((xCnt, yCnt), xyCnt)) => ((x,y), xCnt, yCnt, xyCnt))
    

    另一种选择是使用 map Partition 并简单地处理可迭代对象并跨分区合并分辨率。

    【讨论】:

      【解决方案2】:

      我也是 Spark 的新手,但我知道该怎么做。我不知道这是否是完美的解决方案,但我认为分享它不会有害。

      我会做的可能是 filter() 为值 1 创建一个数据帧和 filter() 为值 0 为第二个数据帧

      你会得到类似的东西

      第一个数据帧

      做 1
      做 1
      FR 1

      在下一步中,我将 groupBy(y)

      所以你会得到第一个数据帧

      做 1 1

      FR 1

      作为 GroupedData https://spark.apache.org/docs/1.4.0/api/java/org/apache/spark/sql/GroupedData.html

      这还有一个 count() 函数,应该计算每组的行数。不幸的是,我现在没有时间自己尝试,但我还是想尝试并提供帮助。

      编辑:如果这有帮助,请告诉我,否则我会删除答案,以便其他人仍然看看这个!

      【讨论】:

      • 感谢您的回答。这个解决方案是我暂时采用的解决方案,但我不确定它是否会推广到 y 的基数为 3 或更大的情况。粗略地说,它包括显式地制作笛卡尔积(同时考虑到y 只能取两个值)。我认为@z-star 提出了一个更通用的答案。但是,请不要删除您的答案,它仍然有效,并且可能对其他一些用户有用,其他贡献者可能会帮助您改进它。
      【解决方案3】:

      最近,我有同样的任务来计算概率,在这里我想分享一下我基于 Spark 的窗口聚合函数的解决方案:

      // data is your DataFrame with two columns [x,y]
      val cooccurrDF: DataFrame = data
        .groupBy(col("x"), col("y"))
        .count()
        .toDF("x", "y", "count-x-y")
      
      val windowX: WindowSpec = Window.partitionBy("x")
      val windowY: WindowSpec = Window.partitionBy("y")
      
      val countsDF: DataFrame = cooccurrDF
        .withColumn("count-x", sum("count-x-y") over windowX)
        .withColumn("count-y", sum("count-x-y") over windowY)
      countsDF.show()
      

      首先,您对两列的所有可能组合进行分组,并使用 count 来获取同时出现的次数。窗口聚合 windowX 和 windowY 允许对聚合行求和,因此您将获得 x 或 y 列的计数。

      +---+---+---------+-------+-------+
      |  x|  y|count-x-y|count-x|count-y|
      +---+---+---------+-------+-------+
      |  0| MK|        1|     17|      1|
      |  0| MX|        1|     17|      1|
      |  1| EE|        1|      3|      1|
      |  0| CN|        1|     17|      1|
      |  1| RU|        1|      3|      2|
      |  0| RU|        1|     17|      2|
      |  0| CL|        1|     17|      1|
      |  0| ES|        1|     17|      1|
      |  0| KR|        1|     17|      1|
      |  0| US|        2|     17|      2|
      |  1| FR|        1|      3|      2|
      |  0| FR|        1|     17|      2|
      |  0| TN|        2|     17|      2|
      |  0| IT|        1|     17|      1|
      |  0| SE|        1|     17|      1|
      |  0| DO|        1|     17|      1|
      |  0| JP|        1|     17|      1|
      |  0| KW|        1|     17|      1|
      +---+---+---------+-------+-------+
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2023-03-03
        • 1970-01-01
        • 1970-01-01
        • 2021-08-05
        • 1970-01-01
        相关资源
        最近更新 更多