Scala RDD按patientID分组获取每组最早日期的实现问题
问题分析
你的代码存在两处核心错误:
groupBy操作返回的RDD元素为(patientID, 同ID所有记录的迭代器)的二元组,不存在date字段,调用x.date会直接编译失败。- 直接对分组后的RDD执行
sortBy+take(1)是取整个RDD排序后的第一条记录,并非每个分组内取最早的记录。
正确实现
首先假设你的case class定义如下:
import org.apache.spark.rdd.RDD case class Medication(patientID: String, date: String) val medication: RDD[Medication] = // 你的原始RDD
方案1:reduceByKey实现(推荐,大数据量场景性能更优)
不需要把同ID的所有记录都拉到同一个节点处理,避免数据倾斜风险:
val out = medication .map(record => (record.patientID, record)) .reduceByKey { (a, b) => // ISO格式的时间字符串字典序和时间顺序一致,可直接比较 if (a.date.compareTo(b.date) < 0) a else b } .values
如果需要更严谨的时间类型校验,可以转换为Instant比较:
import java.time.Instant val out = medication .map(record => (record.patientID, record)) .reduceByKey { (a, b) => if (Instant.parse(a.date).isBefore(Instant.parse(b.date))) a else b } .values
方案2:groupByKey实现(适合小数据量场景)
val out = medication .map(record => (record.patientID, record)) .groupByKey() .mapValues(records => records.minBy(_.date)) .values
输出结果
直接遍历打印即可得到你需要的格式:
println("patientID\tdate") out.collect().foreach(record => println(s"${record.patientID}\t${record.date}"))
内容的提问来源于stack exchange,提问作者user2236600
相关产品推荐
相关产品推荐

