首页
学习
活动
专区
圈层
工具
发布
社区首页 >问答首页 >使用reduceByKey时比较日期

使用reduceByKey时比较日期
EN

Stack Overflow用户
提问于 2016-05-04 19:45:14
回答 2查看 1.3K关注 0票数 1

在scala中,我见过reduceByKey((x: Int , y Int) => x + y),但我想将一个值作为字符串进行迭代,并进行一些比较。我们可以通过reduceByKey(x: String , y: String)使用reduceByKey吗?

代码:

代码语言:javascript
复制
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将值作为字符串进行迭代

EN

回答 2

Stack Overflow用户

发布于 2016-05-04 20:16:21

Spark中的PairRDDFunctions.reduceByKey可以在RDD[(K, V)]形式的任何RDD上工作。reduceByKey将接受类型K (可以使用相等检查对其进行适当比较),并为任何类型的(V, V) => V调用函数V

下面是一个简短的(Int,String)元组示例,其中减少了两个字符串:

代码语言:javascript
复制
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对象。这基本上将代码简化为:

代码语言:javascript
复制
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对象的第一个版本。

票数 1
EN

Stack Overflow用户

发布于 2016-05-04 21:10:29

当我试图将Docsss的类型声明为string时,我没有声明任何类型:

代码语言:javascript
复制
 val docsss : String = rows.map(r => (r(2),r(1)))

很能说明问题

代码语言:javascript
复制
  type mismatch; found : org.apache.spark.rdd.RDD[(Any, Any)]
                 required: String    val rows: RDD[Row] = df.rdd
票数 0
EN
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/37026857

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档