【问题标题】:How to create Spark RDD from an iterator?如何从迭代器创建 Spark RDD?
【发布时间】:2015-09-13 09:22:42
【问题描述】:

为了清楚起见,我不是从像这样的数组/列表中寻找 RDD

List<Integer> list = Arrays.asList(1, 2, 3, 4, 5, 6, 7); // sample
JavaRDD<Integer> rdd = new JavaSparkContext().parallelize(list);

如何在不完全缓冲内存的情况下从 java 迭代器创建 spark RDD?

Iterator<Integer> iterator = Arrays.asList(1, 2, 3, 4).iterator(); //sample iterator for illustration
JavaRDD<Integer> rdd = new JavaSparkContext().what("?", iterator); //the Question

补充问题:

是否需要源可重新读取(或能够多次读取)才能为 RDD 提供弹性?换句话说,既然迭代器本质上是只读的,那么是否有可能从迭代器创建弹性分布式数据集(RDD)?

【问题讨论】:

  • " 没有在内存中完全缓冲" .?你的 Iterator 不是已经在内存中了吗?
  • 数据无论如何都会被加载到内存中。但在我看来,您可以使用 Spark Streaming 来读取输入,因为您的只读一次迭代器可能被视为数据流。
  • @KcDoD 没有。此问题中使用的用于说明。

标签: apache-spark spark-streaming


【解决方案1】:

正如其他人所说,您可以使用 spark 流式处理,但对于纯 spark,您不能,原因是您所要求的内容与 spark 的模型背道而驰。让我解释。 为了分配和并行化工作,spark 必须将其分成块。从 HDFS 读取时,HDFS 为 Spark 完成了“分块”,因为 HDFS 文件是按块组织的。 Spark 通常会为每个块生成一个任务。 现在,迭代器只提供对数据的顺序访问,因此 spark 不可能在不读取内存的情况下将其组织成块

也许可以构建一个具有单个可迭代分区的 RDD,但即便如此,也无法说 Iterable 的实现是否可以发送给工作人员。使用 sc.parallelize() 时,spark 创建了实现serializable 的分区,因此每个分区都可以发送给不同的工作人员。 iterable 可以通过网络连接,或者本地 FS 中的文件,因此除非它们在内存中缓冲,否则它们无法发送给工作人员。

【讨论】:

  • 没错.. 这是一个老问题,但是是的,我通过尝试实现自定义 RDD 解决了这个问题。您所说的非常有道理,因为分区必须可序列化才能获得 RDD。序列化迭代器没有意义。感谢您的确认。
【解决方案2】:

超级老问题,但我只会在序列化后在 flatMap 中创建迭代器。

var ranges = Arrays.asList(Pair.of(1,7), Pair.of(0,5));
JavaRDD<Integer> data = sparkContext.parallelize(ranges).flatMap(pair -> Flux.range(pair.left(), pair.right()).toStream().iterator());

【讨论】:

    猜你喜欢
    • 2016-09-27
    • 2015-05-28
    • 1970-01-01
    • 2015-06-14
    • 1970-01-01
    • 1970-01-01
    • 2018-05-20
    • 2017-06-15
    • 1970-01-01
    相关资源
    最近更新 更多