在scala中,我见过reduceByKey((x: Int , y Int) => x + y),但我想将一个值作为字符串进行迭代,并进行一些比较。我们可以通过reduceByKey(x: String , y: String)使用reduceByKey吗?
代码:
val sparkConf = new SparkConf().setMaster("local").setAppName("Spark AVRO Read")
val sc = new SparkContext(sparkConf)
val inPath= "/home/053764/episodes.avro"
val sqlContext = new SQLContext(sc)
val df = sqlContext.read.avro(inPath)
val rows: RDD[Row] = df.rdd
val doc = df.select("doctor").rdd.map(r => r(0) val docsss = rows.map(r => (r(2), r(1)))
val reduce = docsss.reduceByKey((first, second) => {
val firstDate = LocalDateTime.parse(first)
val secondDate = LocalDateTime.parse(second)
if (firstDate.isBefore(secondDate)) first else second
})请让我知道如何使用reduce by key将值作为字符串进行迭代
发布于 2016-05-04 20:16:21
Spark中的PairRDDFunctions.reduceByKey可以在RDD[(K, V)]形式的任何RDD上工作。reduceByKey将接受类型K (可以使用相等检查对其进行适当比较),并为任何类型的(V, V) => V调用函数V。
下面是一个简短的(Int,String)元组示例,其中减少了两个字符串:
val sc = new SparkContext(conf)
val rdd = sc.parallelize(Seq((1, "01/01/2014"), (1, "02/01/2014")))
rdd.reduceByKey((first, second) => {
val firstDate = LocalDateTime.parse(first)
val secondDate = LocalDateTime.parse(second)
if (firstDate.isBefore(secondDate)) first else second
})编辑:
正如@TheArchetypalPaul正确地指出的那样,由于日期的格式是年/补零月/补零日,因此您可以利用lexicographical order并比较两个String值,而不是将它们解析为DateTime对象。这基本上将代码简化为:
val sc = new SparkContext(conf)
val rdd = sc.parallelize(Seq((1, "01/01/2014"), (1, "02/01/2014")))
rdd.reduceByKey((first, second) => if (first > second) first else second)注意:这确实会限制您使用的特定格式。如果这种情况发生变化,您最好使用为日期创建LocalDateTime对象的第一个版本。
发布于 2016-05-04 21:10:29
当我试图将Docsss的类型声明为string时,我没有声明任何类型:
val docsss : String = rows.map(r => (r(2),r(1)))很能说明问题
type mismatch; found : org.apache.spark.rdd.RDD[(Any, Any)]
required: String val rows: RDD[Row] = df.rddhttps://stackoverflow.com/questions/37026857
复制相似问题