【发布时间】:2015-08-25 12:37:04
【问题描述】:
我有一份工作就这样结束
val iteratorRDD: RDD[Iterator[SomeClass]] = ....
val results = iteratorRDD.map( iterator => iterator.toSeq)
.collect
迭代器是惰性的,即它们在访问其项目时计算数据,这里的 toSeq 基本上会迭代地调用 .next()。
现在,这种计算速度很慢,我想在生成迭代器后立即获取它们的输出,基本上在每个iterator.next() 处。原因是后面的步骤(在本地运行)正在按顺序处理项目:f(all the first items),然后是f(all the seconds),等等......我需要尽快得到这些,因此在工作结束之前。
spark 是否提供某种手段来将中间结果作为某种流检索?或者可能存在迭代器可以向其发送中间数据的分布式数据结构?
我可以做的是设置一个 web 服务来充当这样的缓冲区:它会监听每次调用 iterator.next() 时发送的数据。然后让我的主程序调用该网络服务以获取它存储的内容。但我不喜欢让所有工作人员都与外部服务通信的想法。
【问题讨论】:
标签: apache-spark streaming output