【问题标题】:how to skip empty rdd when join in spark加入spark时如何跳过空rdd
【发布时间】:2015-09-10 02:27:16
【问题描述】:

我想从 Cassandra 获得 2 rdd,然后加入他们。我想跳过空值。

def extractPair(rdd: RDD[CassandraRow]) = {
    rdd.map((row: CassandraRow) => {

     val name = row.getName("name")
     if (name == "")
         None   //join wrong
     else
        (name, row.getUUID("object"))

    })
  }

  val rdd1 = extractPair(cassRdd1)
  val rdd2 = extractPair(cassRdd2)
  val joinRdd = rdd1.join(rdd2)  //"None" join wrong

使用 flatMap 可以解决这个问题,但我想知道如何使用 map 解决这个问题

def extractPair(rdd: RDD[CassandraRow]) = {
        rdd.flatMap((row: CassandraRow) => {

         val name = row.getName("name")
         if (name == "")
             seq()
         else
            Seq((name, row.getUUID("object")))

        })
      }

【问题讨论】:

    标签: join dictionary cassandra apache-spark flatmap


    【解决方案1】:

    仅使用map 是不可能的。您需要使用filter 跟进。但是您仍然最好将有效结果包装在Some 中。但是,那么你仍然会将它包裹在 Some 结果中......需要第二个 map 来解开它。所以,实际上,你最好的选择是这样的:

    def extractPair(rdd: RDD[CassandraRow]) = {
      rdd.flatMap((row: CassandraRow) => {
        val name = row.getName("name")
        if (name == "") None
        else Some((name, row.getUUID("object")))
      })
    }
    

    Option 可隐式转换为可扁平化类型,并更好地传达您的方法消息。

    【讨论】:

    • 这是穿着。值连接不是 org.apache.spark.rdd.RDD[Some[(Any, java.util.UUID)]]的成员
    • 你在使用flatMap吗?它应该去掉Some
    • 我给出的 flatMap 代码可以工作,但我知道如何使用 map。因为我认为 map 比 flatMap 高效。
    • 您认为地图效率更高的原因是什么?我在答案的第一部分解决了您的地图问题……您甚至读过吗?还是只是尝试复制代码?
    • 是的,我已阅读您的建议,但我不明白。 Map用于输入一输出一,FlatMap用于输入一输出多。所以我认为Map效率更高。
    猜你喜欢
    • 2021-09-08
    • 2020-06-25
    • 2019-02-07
    • 2018-11-19
    • 1970-01-01
    • 2019-03-25
    • 2015-06-22
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多