【问题标题】:Spark filter and count big RDD multiple timesSpark过滤并多次计算大RDD
【发布时间】:2016-10-21 10:28:23
【问题描述】:

假设我有一个 RDD[(String, Int)],如下例所示:

(A, 0)
(B, 0)
(C, 1)
(D, 0)
(E, 2)
(F, 1)
(G, 1)
(H, 3)
(I, 2)
(J, 0)
(K, 3)

我想有效地打印包含 0、1、2 等的记录总数。 由于 RDD 包含数百万个条目,因此我希望尽可能高效地执行此操作。

此示例的输出将返回如下内容:

Number of records containing 0 = 4
Number of records containing 1 = 3
Number of records containing 2 = 2
Number of records containing 3 = 2

目前我尝试通过在大 RDD 上执行过滤器来实现这一点,然后分别对 0、1、2、.. 执行 count()。我正在使用 Scala。

有没有更有效的方法来做到这一点?我已经缓存了 RDD,但我的程序仍然内存不足(我已将驱动程序内存设置为 5G)。

编辑: 正如 Tzach 所建议的,我现在使用 countByKey:

rdd.map(_.swap).countByKey()

我是否可以通过将字符串值更改为元组(其中第二个元素是“m”或“f”)来细化这一点,然后获取此元组的第二个元素的每个唯一值的每个键的计数?

例如:

(A,m), 0)
(B,f), 0)
(C,m), 1)
(D,m), 0)
(E,f), 2)
(F,f), 1)
(G,m), 1)
(H,m), 3)
(I,f), 2)
(J,f), 0)
(K,m), 3)

会导致

((0,m), 2)
((0,f), 2)
((1,m), 2)
((1,f), 1)
((2,m), 0)
((2,f), 2)
((3,m), 2)
((3,f), 0)

提前致谢!

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    您可以为此使用方便的countByKey - 只需事先交换输入中的位置以使数值成为键:

    val rdd = sc.parallelize(Seq(
      ("A", 0), ("B", 0), ("C", 1), ("D", 0), ("E", 2),
      ("F", 1), ("G", 1), ("H", 3), ("I", 2), ("J", 0), ("K", 3)
    ))
    
    rdd.map(_.swap).countByKey().foreach(println)
    // (0,4)
    // (1,3)
    // (3,2)
    // (2,2)
    

    编辑countByKey 完全符合它的意思 - 所以无论你想使用什么键,只需将你的 RDD 转换为元组的左侧部分,例如:

    rdd.map { case ((a, b), i) => ((i, b), a) }.countByKey()
    

    或:

    rdd.keyBy { case ((_, b), i) => (i, b) }.countByKey()
    

    【讨论】:

    • 谢谢。我已经编辑了我的答案,以便对问题进行更细化。
    • 相应地编辑了答案(尽管后续问题并不真正符合 SO 礼仪 - 未来的读者更难理解,有时对响应者不公平)
    猜你喜欢
    • 2019-07-08
    • 1970-01-01
    • 2017-11-06
    • 2015-06-15
    • 2019-07-16
    • 2020-12-14
    • 2017-02-20
    • 2015-06-27
    • 1970-01-01
    相关资源
    最近更新 更多