我编写了一个获取DataFrame的类,对其进行了一些计算并可以导出结果.数据帧由密钥列表生成.我知道我现在正以非常低效的方式做到这一点:
var l = List(34, 32, 132, 352) // Scala List
l.foreach{i =>
val data:DataFrame = DataContainer.getDataFrame(i) // get DataFrame
val x = new MyClass(data) // initialize MyClass with new Object
x.setSettings(...)
x.calcSomething()
x.saveResults() // writes the Results into another Dataframe that is saved to HDFS
}
Run Code Online (Sandbox Code Playgroud)
我认为Scala列表上的foreach不是平行的,所以我怎样才能避免在这里使用foreach?DataFrames的计算可以并行进行,因为计算结果不是为下一个DataFrame输入的 - 我该如何实现呢?
非常感谢!!
__编辑:
我试图做的:
val l = List(34, 32, 132, 352) // Scala List
var l_DF:List[DataFrame] = List()
l.foreach{ i =>
DataContainer.getDataFrame(i)::l //append DataFrame to List of Dataframes
}
val rdd:DataFrame …Run Code Online (Sandbox Code Playgroud) 我试图找到一种方法来计算给定数据帧的中位数.
val df = sc.parallelize(Seq(("a",1.0),("a",2.0),("a",3.0),("b",6.0), ("b", 8.0))).toDF("col1", "col2")
+----+----+
|col1|col2|
+----+----+
| a| 1.0|
| a| 2.0|
| a| 3.0|
| b| 6.0|
| b| 8.0|
+----+----+
Run Code Online (Sandbox Code Playgroud)
现在我想做那样的事情:
df.groupBy("col1").agg(calcmedian("col2"))
结果应如下所示:
+----+------+
|col1|median|
+----+------+
| a| 2.0|
| b| 7.0|
+----+------+`
Run Code Online (Sandbox Code Playgroud)
因此calcmedian()必须是UDAF,但问题是,UDAF的"evaluate"方法只需要一行,但我需要整个表来对值进行排序并返回中位数...
// Once all entries for a group are exhausted, spark will evaluate to get the final result
def evaluate(buffer: Row) = {...}
Run Code Online (Sandbox Code Playgroud)
这有可能吗?或者还有另一个不错的解决方法吗?我想强调,我知道如何计算"一组"数据集的中位数.但我不想在"foreach"循环中使用此算法,因为这是低效的!
谢谢!
编辑:
这是我到目前为止所尝试的:
object calcMedian extends UserDefinedAggregateFunction {
// Schema you get as an input
def …Run Code Online (Sandbox Code Playgroud)