I am new to Scala/Spark and I have RDD of case class
case class Info(key1 : String, key2 : String, key3 : String)
I want to transfer RDD[Info] into RDD[JsString] and save it to ElasticSearch, I use play.api.libs and define write converter:
implicit val InfoWrites = new Writes[Info]{
def writes(i : Info): JsObject = Json.obj(
"key1" -> i.key1,
"key2" -> i.key2,
"key3" -> i.key3
)
}
then I define implicit class to use save func:
implicit class Saver(rdd : RDD[Info]) {
def save() : Unit = {
rdd.map{ i => Json.toJson(i).toString }.saveJsonToEs("resource"))
}
}
So I can save RDD[Info] with
infoRDD.save()
But I keep get the "Task not serializable" error with Json.toJson() in rdd.map()
I also try to define serializeable object like this
object jsonUtils extends Serializable{
def toJsString(i : Info) : String = {
Json.toJson(i).toString()
}
}
rdd.map{ i => jsonUtils.toJsString(i) }
but keep getting error "Task not serializable"
How to change the code ? Thank you !