Spark/Scala RDD按user_id分组用最近非空值填充search_id空值实现方案
原生RDD实现方案
核心实现思路
- 先将原始RDD转换为以
user_id为key的键值对RDD - 按
user_id分组后,对每个用户的所有行为数据按timestamp升序排序,保证时间顺序正确 - 遍历每个用户排序后的行为序列,维护一个变量记录当前最近的非空
search_id,遍历过程中:- 若当前行
search_id非空,直接返回当前值,同时更新最近非空search_id - 若当前行
search_id为空,用维护的最近非空search_id填充,无有效值则保留null
- 若当前行
- 最后将处理后的数据展开恢复为原始结构即可
代码实现(Scala)
如果你的原始RDD是文本行类型,先做结构化解析:
val rawRDD = sc.textFile("你的数据源路径") // 过滤表头行 .filter(!_.startsWith("timestamp")) .map(line => { // 移除[]并按逗号分割 val fields = line.replaceAll("\\[|\\]", "").split(",").map(_.trim) val searchId = if (fields(2) == "null") null else fields(2) // 转换为三元组:(时间戳, 用户ID, 搜索ID) (fields(0), fields(1), searchId) })
核心处理逻辑:
val resultRDD = rawRDD // 转换为(user_id, (timestamp, search_id))键值对 .map(item => (item._2, (item._1, item._3))) // 按用户ID分组 .groupByKey() // 对每个用户的记录做填充处理 .flatMapValues(records => { // 先按时间戳升序排序,保证时间顺序正确 val sortedRecords = records.toList.sortBy(_._1) // 维护最近的有效搜索ID var lastValidSearchId: String = null sortedRecords.map { case (ts, searchId) => if (searchId != null) { lastValidSearchId = searchId (ts, searchId) } else { (ts, lastValidSearchId) } } }) // 恢复为原始的三元组结构:(timestamp, user_id, search_id) .map { case (userId, (ts, searchId)) => (ts, userId, searchId) }
性能优化说明
如果存在数据量极大、部分用户行为记录特别多的场景,为了避免groupByKey带来的数据倾斜问题,可以采用二次排序优化方案:
- 将
(user_id, timestamp)作为复合key,先对全量RDD按复合key排序 - 再用
mapPartitions按user_id分组处理,同一个用户的记录会连续出现在同一个分区,避免shuffle时单key数据量过大的问题
内容的提问来源于stack exchange,提问作者justino666
相关产品推荐
相关产品推荐

