【发布时间】:2017-09-29 14:15:33
【问题描述】:
问题出在 map 函数中,同时进行案例类提取。案例类不可序列化。我已经隐式定义了格式DefaultFormats。
package org.apache.flink.quickstart
import java.util.Properties
import com.fasterxml.jackson.databind.{JsonNode, ObjectMapper}
import com.fasterxml.jackson.module.scala.DefaultScalaModule
import org.apache.flink.api.scala._
import org.apache.flink.runtime.state.filesystem.FsStateBackend
import org.apache.flink.streaming.api.scala.{DataStream, StreamExecutionEnvironment}
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer09
import org.apache.flink.streaming.util.serialization.SimpleStringSchema
import org.json4s.DefaultFormats
import org.json4s._
import org.json4s.native.JsonMethods
import scala.util.Try
case class CC(key:String)
object WordCount{
def main(args: Array[String]) {
implicit val formats = org.json4s.DefaultFormats
// kafka properties
val properties = new Properties()
properties.setProperty("bootstrap.servers", "***.**.*.***:9093")
properties.setProperty("zookeeper.connect", "***.**.*.***:2181")
properties.setProperty("group.id", "afs")
properties.setProperty("auto.offset.reset", "earliest")
val env = StreamExecutionEnvironment.getExecutionEnvironment
val st = env
.addSource(new FlinkKafkaConsumer09("new", new SimpleStringSchema() , properties))
.flatMap(raw => JsonMethods.parse(raw).toOption)
// .map(_.extract[CC])
val l = st.map(_.extract[CC])
st.print()
env.execute()
}
}
错误:
INFO [main] (TypeExtractor.java:1804) - 未检测到类的字段 org.json4s.JsonAST$JValue。不能用作 PojoType。将会 作为 GenericType 处理 线程“主”org.apache.flink.api.common.InvalidProgramException 中的异常:任务不是 可序列化 在 org.apache.flink.api.scala.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:172) 在 org.apache.flink.api.scala.ClosureCleaner$.clean(ClosureCleaner.scala:164) 在 org.apache.flink.streaming.api.scala.StreamExecutionEnvironment.scalaClean(StreamExecutionEnvironment.scala:666) 在 org.apache.flink.streaming.api.scala.DataStream.clean(DataStream.scala:994) 在 org.apache.flink.streaming.api.scala.DataStream.map(DataStream.scala:519) 在 org.apache.flink.quickstart.WordCount$.main(WordCount.scala:38) 在 org.apache.flink.quickstart.WordCount.main(WordCount.scala) 引起:java.io.NotSerializableException: org.json4s.DefaultFormats$$anon$4 在 java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1184) 在 java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1548) 在 java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1509) 在 java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1432) 在 java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1178) 在 java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1548) 在 java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1509) 在 java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1432) 在 java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1178) 在 java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:348) 在 org.apache.flink.util.InstantiationUtil.serializeObject(InstantiationUtil.java:317) 在 org.apache.flink.api.scala.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:170) ... 6 更多
Process finished with exit code 1
【问题讨论】:
标签: json scala serialization apache-flink json4s