You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Scala RDD按patientID分组获取每组最早日期的实现问题

问题分析

你的代码存在两处核心错误:

  1. groupBy操作返回的RDD元素为(patientID, 同ID所有记录的迭代器)的二元组,不存在date字段,调用x.date会直接编译失败。
  2. 直接对分组后的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.03 09:18:03