【问题标题】:Spark, efficient way to get a single value from each partitions? accumulator?Spark,从每个分区获取单个值的有效方法?累加器?
【发布时间】: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 的评估(一次在代码中,另一次在后面),这在我的应用程序中很重。而且,cachepersist 会有所帮助,但我无法在我的应用程序中使用。

有什么好的方法可以从每个分区中收集对象吗? 蓄能器呢?

【问题讨论】:

  • 累加器无济于事,因为您不能真正在转换或动作中使用累加值。
  • 你试过select
  • 为什么不能使用缓存或持久化?
  • @DennisJaheruddin 因为df 可能是巨大的,比我的记忆大一百倍。

标签: apache-spark rdd lazy-evaluation partition accumulator


【解决方案1】:

使传递给mapPartitions 的函数返回一个仅包含摘要对象的迭代器。然后就可以直接采集了,不需要额外的过滤。

【讨论】:

  • 你的意思是``` val summaries = df.mapPartitions(_ => /*add Summary Object*/ ).collect() val data = df.mapPartitions(_ => /*数据处理* /)```我觉得解决不了问题,因为还需要两次求值。
【解决方案2】:

为什么不能使用缓存或持久化? ——丹尼斯·贾赫鲁丁

@DennisJaheruddin 因为 df 可能很大,比我的记忆大一百倍。 – 用户2037661

如果要缓存不适合内存的数据帧,可以使用存储级别 MEMORY_AND_DISK。目前调用cachepersist时默认为该存储级别。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2016-05-13
    • 2019-12-28
    • 2018-02-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-02-11
    • 1970-01-01
    相关资源
    最近更新 更多