【发布时间】:2017-09-13 08:10:06
【问题描述】:
我是 Spark 的相对初学者。我有一个宽数据框(1000 列),我想根据相应列是否缺少值来添加列
所以
+----+ |一个 | +----+ | 1 | +----+ |空| +----+ | 3 | +----+变成
+----+-------+ |一个 |管理信息系统 | +----+-------+ | 1 | 0 | +----+-------+ |空| 1 | +----+-------+ | 3 | 1 | +----+-------+这是自定义 ml 转换器的一部分,但算法应该是清晰的。
override def transform(dataset: org.apache.spark.sql.Dataset[_]): org.apache.spark.sql.DataFrame = {
var ds = dataset
dataset.columns.foreach(c => {
if (dataset.filter(col(c).isNull).count() > 0) {
ds = ds.withColumn(c + "_MIS", when(col(c).isNull, 1).otherwise(0))
}
})
ds.toDF()
}
循环遍历列,如果 > 0 个空值创建一个新列。
传入的数据集被缓存(使用 .cache 方法),相关配置设置为默认值。 这目前在一台笔记本电脑上运行,即使只有最少的行,1000 列的运行时间约为 40 分钟。 我认为问题出在数据库中,所以我尝试使用镶木地板文件,但结果相同。查看作业 UI,它似乎正在执行文件扫描以进行计数。
有没有办法改进这个算法以获得更好的性能,或者以某种方式调整缓存?增加 spark.sql.inMemoryColumnarStorage.batchSize 只是让我出现 OOM 错误。
【问题讨论】:
标签: scala apache-spark spark-dataframe