【问题标题】:Spark Closures with Array [duplicate]带数组的 Spark 闭包 [重复]
【发布时间】:2015-08-07 06:14:13
【问题描述】:
我有一个数组,当它在闭包内(它有一些值)但在循环外,数组大小为 0。我想知道是什么导致了这种行为?
我需要在外部可以访问 hArr 以进行批量 HBase 放置。
val hArr = new ArrayBuffer[Put]()
rdd.foreach(row => {
val hConf = HBaseConfiguration.create()
val hTable = new HTable(hConf, tablename)
val hRow = new Put(Bytes.toBytes(row._1.toString))
hRow.add(...)
hArr += hRow
println("hArr: " + hArr.toArray.mkString(","))
})
println("hArr.size: " + hArr.size)
【问题讨论】:
标签:
scala
hadoop
apache-spark
closures
rdd
【解决方案1】:
问题是 rdd 闭包中的任何项目都被复制并使用本地版本。 foreach 只能用于保存到磁盘或类似的东西。
如果你想把它放在一个数组中,那么你可以map 然后collect
rdd.map(row=> {
val hConf = HBaseConfiguration.create()
val hTable = new HTable(hConf, tablename)
val hRow = new Put(Bytes.toBytes(row._1.toString))
hRow.add(...)
hRow
}).collect()
【解决方案2】:
我发现相当多的 Spark 新用户对 mapper 和 reducer 函数如何运行以及它们如何与驱动程序中定义的事物相关联感到困惑。通常,您通过 map 或 foreach 或 reduceByKey 或许多其他变体定义和注册的所有 mapper/reducer 函数都不会在您的驱动程序程序中执行。在您的驱动程序中,您只需为 Spark 注册它们以远程和分布式运行它们。当这些函数引用您在驱动程序中实例化的某些对象时,您实际上创建了一个“闭包”,它在大多数情况下都可以编译。但通常这不是您想要的,您通常会在运行时遇到问题,通过看到 NotSerializable 或 ClassNotFound 异常。
您可以通过 foreach() 变体远程完成所有输出工作,或者尝试通过调用 collect() 将所有数据收集回驱动程序以进行输出。但是要小心 collect() ,因为它会将所有数据从分布式节点收集到驱动程序的程序中。只有当您完全确定最终汇总的数据很小时,您才会这样做。