【问题标题】:Cannot write/save data to Ignite directly from a Spark RDD无法直接从 Spark RDD 将数据写入/保存到 Ignite
【发布时间】:2018-04-15 16:02:13
【问题描述】:

我尝试使用 jdbc 编写数据帧来点燃,

Spark 版本是:2.1

点燃版本:2.3

JDK:1.8

斯卡拉:2.11.8

这是我的代码 sn-p:

def WriteToIgnite(hiveDF:DataFrame,targetTable:String):Unit = {

  val conn = DataSource.conn
  var psmt:PreparedStatement = null

  try {
    OperationIgniteUtil.deleteIgniteData(conn,targetTable)

    hiveDF.foreachPartition({
      partitionOfRecords => {
        partitionOfRecords.foreach(
          row => for ( i <- 0 until row.length ) {
            psmt = OperationIgniteUtil.getInsertStatement(conn, targetTable, hiveDF.schema)
            psmt.setObject(i+1, row.get(i))
            psmt.execute()
          }
        )
      }
    })

  }catch {
    case e: Exception =>  e.printStackTrace()
  } finally {
    conn.close
  }
}

然后我在 spark 上运行,它打印错误消息:

org.apache.spark.SparkException:任务不可序列化 在 org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:298) 在 org.apache.spark.util.ClosureCleaner$.org$apache$spark$util$ClosureCleaner$$clean(ClosureCleaner.scala:288) 在 org.apache.spark.util.ClosureCleaner$.clean(ClosureCleaner.scala:108) 在 org.apache.spark.SparkContext.clean(SparkContext.scala:2094) 在 org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1.apply(RDD.scala:924) 在 org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1.apply(RDD.scala:923) 在 org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151) 在 org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112) 在 org.apache.spark.rdd.RDD.withScope(RDD.scala:362) 在 org.apache.spark.rdd.RDD.foreachPartition(RDD.scala:923) 在 org.apache.spark.sql.Dataset$$anonfun$foreachPartition$1.apply$mcV$sp(Dataset.scala:2305) 在 org.apache.spark.sql.Dataset$$anonfun$foreachPartition$1.apply(Dataset.scala:2305) 在 org.apache.spark.sql.Dataset$$anonfun$foreachPartition$1.apply(Dataset.scala:2305) 在 org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:57) 在 org.apache.spark.sql.Dataset.withNewExecutionId(Dataset.scala:27​​65) 在 org.apache.spark.sql.Dataset.foreachPartition(Dataset.scala:2304) 在 com.pingan.pilot.ignite.common.OperationIgniteUtil$.WriteToIgnite(OperationIgniteUtil.scala:72) 在 com.pingan.pilot.ignite.etl.HdfsToIgnite$.main(HdfsToIgnite.scala:36) 在 com.pingan.pilot.ignite.etl.HdfsToIgnite.main(HdfsToIgnite.scala) 在 sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) 在 sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) 在 sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) 在 java.lang.reflect.Method.invoke(Method.java:498) 在 org.apache.spark.deploy.SparkSubmit$.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:738) 在 org.apache.spark.deploy.SparkSubmit$.doRunMain$1(SparkSubmit.scala:187) 在 org.apache.spark.deploy.SparkSubmit$.submit(SparkSubmit.scala:212) 在 org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:126) 在 org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala) 引起:java.io.NotSerializableException: org.apache.ignite.internal.jdbc2.JdbcConnection 序列化栈: - 对象不可序列化(类:org.apache.ignite.internal.jdbc2.JdbcConnection,值: org.apache.ignite.internal.jdbc2.JdbcConnection@7ebc2975) - 字段(类:com.pingan.pilot.ignite.common.OperationIgniteUtil$$anonfun$WriteToIgnite$1, 名称:conn$1,类型:接口 java.sql.Connection) - 对象(com.pingan.pilot.ignite.common.OperationIgniteUtil$$anonfun$WriteToIgnite$1 类, ) 在 org.apache.spark.serializer.SerializationDebugger$.improveException(SerializationDebugger.scala:40) 在 org.apache.spark.serializer.JavaSerializationStream.writeObject(JavaSerializer.scala:46) 在 org.apache.spark.serializer.JavaSerializerInstance.serialize(JavaSerializer.scala:100) 在 org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:295) ... 27 更多

有人知道我可以修复它吗? 谢谢!

【问题讨论】:

标签: java scala apache-spark jdbc ignite


【解决方案1】:

这里的问题是您无法序列化与 Ignite DataSource.conn 的连接。您提供给 forEachPartition 的闭包包含连接作为其范围的一部分,这就是 Spark 无法序列化它的原因。

幸运的是,Ignite 提供了 RDD 的自定义实现,允许您将值保存到它。您需要先创建一个IgniteContext,然后检索Ignite 的共享RDD,它提供对Ignite 的分布式访问以保存您的RDD 的Row

val igniteContext = new IgniteContext(sparkContext, () => new IgniteConfiguration())
...

// Retrieve Ignite's shared RDD
val igniteRdd = igniteContext.fromCache("partitioned")
igniteRDD.saveValues(hiveDF.toRDD)

更多信息请访问Apache Ignite documentation

【讨论】:

    【解决方案2】:

    你必须扩展 Serializable 接口。

    object Test extends Serializable { 
      def WriteToIgnite(hiveDF:DataFrame,targetTable:String):Unit = {
       ???
      }
    }
    

    我希望它能解决你的问题。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-03-25
      • 2023-03-14
      • 2016-01-30
      • 1970-01-01
      • 1970-01-01
      • 2017-03-26
      相关资源
      最近更新 更多