【发布时间】:2017-07-17 09:00:14
【问题描述】:
我有以下目录结构:
/数据/模型A
/数据/模型B
/数据/模型C ..
这些文件中的每一个都有格式(id、score)的数据,我必须分别为它们做以下操作-
1) 按分数分组并按降序排列分数(DF_1: score,count)
2) 从 DF_1 计算每个排序的分数组的累积频率(DF_2:score, count, cumFreq)
3) 从 DF_2 中选择介于 5-10 之间的累积频率 (DF_3: score, cumFreq)
4) 从 DF_3 中选择最低分数(DF_4: score)
5) 从文件中选择所有分数大于 DF_4 分数的 id 并保存
我可以通过将目录读取为 wholeTextFile 并为所有模型创建一个公共数据框,然后在模型上使用 group by 来做到这一点。
我想做-
val scores_file = sc.wholeTextFiles("/data/*/")
val scores = scores_file.map{ line =>
//step 1
//step 2
//step 3
//step 4
//step 5 : save as line._1
}
这将有助于分别处理每个文件,并避免分组。
【问题讨论】:
标签: scala apache-spark dataframe apache-spark-sql rdd