【问题标题】:collect RDD output as a stream将 RDD 输出收集为流
【发布时间】: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


    【解决方案1】:

    我这样做没有任何意义。当您不想在本地内存中创建副本时,一个迭代器可以很好地遍历,但 Spark 的工作方式不同。您的迭代器分布在多个执行器中(在单独的节点中,具有单独的内存),因此当您调用 collect 时,您将强制它们被迭代并发送到主服务器,它们将被加载到内存中。您根本无法从 master 对 executor 中的数据进行惰性评估。

    您应该努力将计算发送给数据,而不是相反,特别是如果您要在每个序列上运行相同的代码!例如:

    val results = iteratorRDD
      .map(iter => f(iter)) // Whatever f() returns.
      .collect()
    

    然后,您在执行器上懒惰地并行评估您的迭代器,只将实际结果带给主节点。

    【讨论】:

    • 我没有谈论输入数据,所以认为它在迭代器中:val inputData = sc.textFile(hdfs_url) ; iteratorRdd = inputData.map( line => makeIterator(line) 因此迭代器是围绕数据进行的计算。现在我的观点是它(懒惰地)在每个工作人员上计算一系列值。我想在所有计算之前收到第一个。
    • 是的,我明白了。所以让我重新表述这个问题:为什么你必须在主人身上接收这些?为什么要收集它们来运行“后续步骤”?
    • 后面的步骤不能在公有云上运行(要不然我找到了隐藏敏感信息的方法,但是这很难)而且输出大小没有那么大。
    • 那么恐怕没有.collectIterator() 功能可以满足您的需求,考虑到Spark 的工作原理,它们可能永远不会。
    猜你喜欢
    • 2021-01-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-08-17
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多