【发布时间】:2020-01-14 19:38:27
【问题描述】:
我在驱动程序中创建了 spark jdbc 单例连接,并计划在执行程序中使用连接。我得到以下异常。 org.apache.spark.SparkException:任务不可序列化
Spark主类内部:
object ExecutorConnection {
private var connection: Connection = null
val url = prop.getProperty("url")
val user = prop.getProperty("user")
val pwd = prop.getProperty("password")
val driver = prop.getProperty("driver")
Class.forName(driver)
def getConnection(url: String, username: String, password: String): Connection = synchronized {
if (connection == null) {
connection = DriverManager.getConnection(url, username, password)
Class.forName(driver)
connection.setAutoCommit(false)
}
connection
}
lazy val createConnection = getConnection(url, user, pwd)
}
我有多个具有不同架构的数据帧(df1,df2,df3),我计划在驱动程序级别创建连接并序列化连接并将其用于所有数据帧。
df1.rdd.repartition(2).mapPartitions((d) => Iterator(d)).foreach { partition =>
val conn = ExecutorConnection.createConnection
var ps: PreparedStatement = null
partition.grouped(1).foreach(batch => {
batch.foreach { x =>
{
ps = conn.prepareStatement(SqlString)
ps.addBatch()
conn.commit()
}
}
})
}
【问题讨论】:
-
我认为你应该在 mapPartitions 之后调用 foreach 分区而不是 foreach。基本上,您的 mapPartitions 什么都不做,然后您正在浏览记录,而不是分区。您不能序列化连接,只能在工作人员处从头开始创建它。
标签: scala apache-spark apache-spark-sql