【发布时间】:2019-05-17 17:01:12
【问题描述】:
我想要实现的是,对于以下 DataFrame:
-------------------------
| FOO | BAR | BAZ |
| lorem | ipsum | dolor |
| sit | amet | dolor |
| lorem | lorem | dolor |
-------------------------
生成以下输出:
Map(
FOO -> List("lorem", "sit"),
BAR -> List("ipsum", "amet", "lorem"),
BAZ -> List("dolor")
)
这是我想出的 Scala 代码:
val df = data.distinct
df.columns.map((key) => {
val distinctValues = df
.select(col(key))
.collect
.map(df => df.getString(0))
.toList
.distinct
(key, distinctValues)
}).toMap
我已经尝试过使用 RDD 来替代此代码,不知何故,它们的速度提高了大约 30%,但问题仍然存在: 这一切都非常低效。
我在本地 Cassandra 上运行 Spark,该 Cassandra 托管只有 1000 行的示例数据集,但这些操作会生成大量日志,并且需要 7 秒以上才能完成。
我是不是做错了什么,有更好的方法吗?
【问题讨论】:
-
df.select(df.columns map (c => collect_set(c) as c): _*).first.getValuesMap[Seq[String]](df.columns)将是一个小的改进,但总体思路是不可扩展的,并且在通用集群上的 Spark 中通常会出现第二个延迟。 -
我认为这可能是一个重复的问题 - stackoverflow.com/questions/37949494/…
-
@user6910411 确实提高了很多性能,谢谢!您能否详细说明为什么这是不可扩展的?因为输出会变得太大?
-
好吧。您收集所有唯一值,并转换为本地结构,以便当且仅当最终结果足够小以由每个节点(驱动程序和执行程序)在内存中处理时它才能工作。如果您先验地知道基数很小,那么可以工作,但是假设有人在那里放了一个随机实数:) 如果您愿意放宽您的要求,您可以melt 框架,然后只取不同的值来保持分布。跨度>
标签: scala apache-spark dataframe apache-spark-sql rdd