【问题标题】:How map Key/Value pairs between two separate RDDs?如何在两个单独的 RDD 之间映射键/值对?
【发布时间】:2017-12-06 06:50:09
【问题描述】:

仍然是 Scala 和 Spark 的初学者,我想我只是在这里无脑。我有两个 RDD,其中一种是:-

((String, String), Int) = ((" v67430612_serv78i"," fb_201906266952256"),1)

其他类型:-

(String, String, String) = (r316079113_serv60i,fb_100007609418328,-795000)

可以看出,两个RDD的前两列格式相同。基本上它们是ID,一个是'tid',另一个是'uid'。

问题是这样的:

有没有一种方法可以比较两个 RDD,使得 tid 和 uid 在两者中都匹配,并且相同匹配 id 的所有数据都显示在一行中而没有任何重复?

例如:如果我得到两个 RDD 之间的 tid 和 uid 匹配

((String, String), Int) = ((" v67430612_serv78i"," fb_201906266952256"),1)

(String, String, String) = (" v67430612_serv78i"," fb_201906266952256",-795000)

那么输出是:-

((" v67430612_serv78i"," fb_201906266952256",-795000),1)

两个 RDD 中的 ID 没有任何固定的顺序。它们是随机的,即相同的 uid 和 tid 序列号在两个 RDD 中可能不对应。

另外,如果第一个 RDD 类型保持不变但第二个 RDD 更改为类型,解决方案将如何变化:-

((String, String, String), Int) = ((daily_reward_android_5.76,fb_193055751144610,81000),1)

我必须在不使用 Spark SQL 的情况下执行此操作。

【问题讨论】:

    标签: scala apache-spark string-matching


    【解决方案1】:

    我建议您将rdds 转换为dataframes 并申请join 以方便使用。

    你的第一个dataframe应该是

    +------------------+-------------------+-----+
    |tid               |uid                |count|
    +------------------+-------------------+-----+
    | v67430612_serv78i| fb_201906266952256|1    |
    +------------------+-------------------+-----+
    

    第二个dataframe应该是

    +------------------+-------------------+-------+
    |tid               |uid                |amount |
    +------------------+-------------------+-------+
    | v67430612_serv78i| fb_201906266952256|-795000|
    +------------------+-------------------+-------+
    

    那么得到最后的输出就是innerjoin as

    df2.join(df1, Seq("tid", "uid"))
    

    输出为

    +------------------+-------------------+-------+-----+
    |tid               |uid                |amount |count|
    +------------------+-------------------+-------+-----+
    | v67430612_serv78i| fb_201906266952256|-795000|1    |
    +------------------+-------------------+-------+-----+
    

    已编辑

    如果您想在没有 dataframe/spark sql 的情况下执行此操作,那么也可以通过 rdd 方式加入,但您必须进行如下修改

    rdd2.map(x => ((x._1, x._2), x._3)).join(rdd1).map(y => ((y._1._1, y._1._2, y._2._1), y._2._2)) 
    

    仅当您的问题中将rdd1rdd2 分别定义为((" v67430612_serv78i"," fb_201906266952256"),1)(" v67430612_serv78i"," fb_201906266952256",-795000) 时,这才有效。 你应该有最终输出为

    (( v67430612_serv78i, fb_201906266952256,-795000),1)
    

    确保修剪空格的值。这将帮助您确保两个 rdd 在加入时具有相同的 key 值,否则您可能会得到一个空结果。

    【讨论】:

    • 为什么要接受?您需要没有 spark SQL 的帮助,不是吗?
    • 没有冒犯!官方是的。但非正式地,这也教会了我一些新的东西,因为我以前从未单独使用 SQL 或使用过 Spark。这就是接受和支持的原因。
    • 我认为您可以使用 RDD map、mapPartitions 和 join 方法来实现您的目标。
    • @QueepyDev 人们将停止调查您的问题,因为您已经接受了这个问题的答案。我同意你可以投票,因为它对你有用。
    • 但是RDD类型不同。这不会导致类型不匹配吗? @汤姆
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-06-07
    • 2017-10-19
    • 1970-01-01
    • 2015-03-24
    • 1970-01-01
    • 2013-03-20
    相关资源
    最近更新 更多