【发布时间】:2016-12-15 17:19:21
【问题描述】:
我在 spark 中有一个操作,应该对数据框中的几列执行。一般来说,有两种可能来指定这样的操作
- 硬编码
handleBias("bar", df)
.join(handleBias("baz", df), df.columns)
.drop(columnsToDrop: _*).show
- 从列名列表动态生成它们
var isFirst = true
var res = df
for (col <- columnsToDrop ++ columnsToCode) {
if (isFirst) {
res = handleBias(col, res)
isFirst = false
} else {
res = handleBias(col, res)
}
}
res.drop(columnsToDrop: _*).show
问题在于动态生成的 DAG 是不同的,当使用更多列时,动态解决方案的运行时间会比硬编码操作增加得更多。
我很好奇如何将优雅的动态构造与快速执行时间结合起来。
对于大约 80 列,这为硬编码变体生成了一个相当不错的图表 对于动态构造的查询,还有一个非常大、可能不太可并行化且速度较慢的 DAG。
当前版本的 spark (2.0.2) 与 DataFrames 和 spark-sql 一起使用
完成最小示例的代码:
def handleBias(col: String, df: DataFrame, target: String = "FOO"): DataFrame = {
val pre1_1 = df
.filter(df(target) === 1)
.groupBy(col, target)
.agg((count("*") / df.filter(df(target) === 1).count).alias("pre_" + col))
.drop(target)
val pre2_1 = df
.groupBy(col)
.agg(mean(target).alias("pre2_" + col))
df
.join(pre1_1, Seq(col), "left")
.join(pre2_1, Seq(col), "left")
.na.fill(0)
}
编辑
使用foldleft 运行任务会生成线性 DAG
并对所有列的函数进行硬编码会导致
两者都比我原来的 DAG 好很多,但硬编码的变体对我来说看起来更好。在 spark 中连接 SQL 语句的字符串可以让我动态生成硬编码的执行图,但这看起来相当难看。您还有其他选择吗?
【问题讨论】:
-
我认为问题在于您的“handleBias”函数非常复杂,您需要为多个列运行它。即使你对许多列进行硬编码,你的 DAG 也会很大,所以问题可能不是“动态”应用,而是应用于许多列。因此,如果您能想出一种方法来调整您的函数以同时处理多个列,那可能会有很大帮助。
-
@DanieldePaula 你有什么方法可以用更简单的方式来表达这种方法,从而减少所需的计算能力?
-
很遗憾,我现在没有太多时间考虑,对不起。如果到明天你还没有找到解决方案,我会看看它。
-
@DanieldePaula 到目前为止我还想不出一个简化的方法。想知道是否可以改进缓存?目前,我在更大的数据集上调用此函数之前使用缓存。
-
@DanieldePaula 你认为我可以通过“连接”列来摆脱其中的一些连接吗? stackoverflow.com/questions/32882529/… 因为在这里我需要一个 concat 操作(可以并行执行)。
标签: apache-spark apache-spark-sql apache-spark-dataset