【发布时间】: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