|
|
1
5
下面是在pyspark中执行相同操作的说明。我同意拉姆的解释。 collectAsMap仅适用于pairedrdd,因此您需要先将数据帧转换为pairedrdd,然后使用collectAsMap函数将其转换为一些字典。 例如,我有一个下面的数据框:
将其转换为键值对rdd
最后,您可以使用collectAsMap将键值对rdd转换为dict
|
|
2
2
首先,我在python/pyspark方面不好,所以我演示了使用scala。。。
你的
所以你不能把它收集成地图。除非你必须做一个
下面是一个完整的例子来证明这一点。 import org.apache.log4j.{Level, Logger}
import org.apache.spark.internal.Logging
import org.apache.spark.sql.SparkSession
/** *
* collectAsMap is only applicable to pairedrdd if you want to do a map then you can do a rdd key by and proceed
*
* @author : Ram Ghadiyaram
*/
object PairedRDDPlay extends Logging {
Logger.getLogger("org").setLevel(Level.OFF)
// Logger.getLogger("akka").setLevel(Level.OFF)
def main(args: Array[String]): Unit = {
val appName = if (args.length > 0) args(0) else this.getClass.getName
val spark: SparkSession = SparkSession.builder
.config("spark.master", "local") //.config("spark.eventLog.enabled", "true")
.appName(appName)
.getOrCreate()
import spark.implicits._
val pairs = spark.sparkContext.parallelize(Array((1, 1,3), (1, 2,3), (1, 3,3), (1, 1,3), (2, 1,3))).toDF("mycol1", "mycol2","mycol3")
pairs.show()
val keyedBy = pairs.rdd.keyBy(_.getAs[Int]("mycol1"))
keyedBy.foreach(x => println("using keyBy-->>" + x))
val myMap = keyedBy.collectAsMap()
println(myMap.toString())
assert(myMap.size == 2)
// val myMap1 = pairs.rdd.collectAsMap()
// println(myMap1.toString())
// assert(myMap1.size == 2)
//Error:(28, 28) value collectAsMap is not a member of org.apache.spark.rdd.RDD[org.apache.spark.sql.Row]
// val myMap1 = pairs.rdd.collectAsMap()
}
}
结果:
答:不,您可以在示例中看到具有多列(即>2)的示例。但你需要把它转换成pairrdd。 |