【问题标题】:what are other approaches instead of using collect() in spark scala除了在 spark scala 中使用 collect() 之外,还有什么其他方法
【发布时间】:2016-09-20 11:27:45
【问题描述】:
我的一段 scala 代码看起来像,
val orgIncInactive = orgIncLatest.filter("(LD_TMST != '' and LD_TMST is not null)").select("ORG_ID").rdd
orgIncInactive.collect.foreach(p => DenormalizedTablesMethodsUtil.hbaseTablePurge(p(0).toString, tableName, connection))
有什么方法可以避免在这里使用 collect() 吗?
我尝试了各种可能性,但最终出现了可序列化的错误。
谢谢。
【问题讨论】:
标签:
scala
hadoop
apache-spark
rdd
【解决方案1】:
取决于您要执行的操作,以及最终导致序列化错误的原因。看起来您正试图将某种数据库连接传递给匿名函数。由于几个原因,这通常会失败。即使您使连接对象本身可序列化(例如通过对对象进行子类化并实现Serializable),您也不能在驱动程序和执行程序之间共享数据库连接。
相反,您需要做的是在每个执行器上创建连接对象,然后使用本地连接对象而不是驱动程序中定义的连接对象。有几种方法可以做到这一点。
一种是使用mapPartitions,它允许您在逻辑运行之前在本地实例化对象。请参阅 here 了解更多信息。
另一种可能性是创建一个在初始化时将连接对象设置为null 或None 的单例对象。然后,您将在对象中定义一个方法,例如“getConnection”,用于检查连接是否已初始化。如果不是,它会初始化连接。然后无论哪种方式,它都会返回有效的连接。
我使用第二种方法比第一种更多,因为它将初始化限制为每个执行程序仅一次,而不是强制每个分区发生一次。