【发布时间】:2019-11-15 08:17:15
【问题描述】:
在我的 spark 应用程序中,每个分区都会生成一个 Object,它很小,并且包含分区中的数据摘要。现在我通过将它们放入也包含主要数据的 Datafram 来收集它们。
val df: DataFrame[(String, Any)] = df.mapPartitions(_ => /*add Summary Object*/ )
val summaries = df.filter(_._1 == "summary").map(_._2).collect()
val data = df.filter(_._1 == "data").map(_._2) // used to further RDD processing
Summary 对象立即使用,data 将用于 RDD 处理。
问题是,代码产生了两次df 的评估(一次在代码中,另一次在后面),这在我的应用程序中很重。而且,cache 或 persist 会有所帮助,但我无法在我的应用程序中使用。
有什么好的方法可以从每个分区中收集对象吗? 蓄能器呢?
【问题讨论】:
-
累加器无济于事,因为您不能真正在转换或动作中使用累加值。
-
你试过
select -
为什么不能使用缓存或持久化?
-
@DennisJaheruddin 因为
df可能是巨大的,比我的记忆大一百倍。
标签: apache-spark rdd lazy-evaluation partition accumulator