Scala中用Apache Flink读取Kafka JSON消息遇序列化错误求助
解决Flink+Scala读取Kafka JSON时的序列化错误
问题场景
尝试用Apache Flink和Scala读取Kafka Topic中的JSON消息,示例消息:
{"brand": "apple", "category": "phone", "productid": "2", "productname":"iphone"}
原有代码如下:
object ProductModel extends Serializable { def readElement(jsonElement: JsValue): productJs= { val brand = (jsonElement \ "brand").as[String] val category = (jsonElement \ "category").as[String] val productid = (jsonElement \ "productid").as[String] val productname = (jsonElement \ "productname").as[String] productJs(brand, category, productid, productname) } case class productJs(brand: String, category: String, productid: String, productname: String) } val env = StreamExecutionEnvironment.getExecutionEnvironment val properties = new Properties() properties.setProperty("bootstrap.servers", "localhost:9092") properties.setProperty("group.id", "flink-consumer-group") properties.setProperty("auto.offset.reset", "earliest") val consumer = new FlinkKafkaConsumer[String]("my-topic", new SimpleStringSchema(), properties) val stream1 = env.addSource(consumer) val stream2 = stream1.map{ x => ProductModel.readElement(Json.parse(x))} stream2.print() env.execute("FlinkKafkaScala")
运行时抛出错误:
org.apache.flink.api.common.InvalidProgramException: Task not serializable
错误原因
Flink的算子任务需要可序列化,才能分发到集群节点执行。此处问题根源:
productJscase class定义在ProductModelobject内部,Scala中嵌套在object里的case class,其序列化会依赖外部object。即便ProductModel继承了Serializable,Flink的序列化机制仍可能因这种嵌套结构触发序列化失败。map算子中引用ProductModel.readElement方法时,Flink序列化匿名函数会尝试序列化ProductModel对象,嵌套case class的序列化逻辑存在隐藏的非序列化依赖。
修正后的代码
方案1:将case class移至顶层作用域
把productJs移到ProductModel外部,消除嵌套依赖,同时遵循Scala命名规范将类名改为首字母大写:
// 顶层case class,独立于object,避免序列化依赖 case class ProductJs(brand: String, category: String, productid: String, productname: String) object ProductModel { def readElement(jsonElement: JsValue): ProductJs = { val brand = (jsonElement \ "brand").as[String] val category = (jsonElement \ "category").as[String] val productid = (jsonElement \ "productid").as[String] val productname = (jsonElement \ "productname").as[String] ProductJs(brand, category, productid, productname) } } val env = StreamExecutionEnvironment.getExecutionEnvironment val properties = new Properties() properties.setProperty("bootstrap.servers", "localhost:9092") properties.setProperty("group.id", "flink-consumer-group") properties.setProperty("auto.offset.reset", "earliest") val consumer = new FlinkKafkaConsumer[String]("my-topic", new SimpleStringSchema(), properties) val stream1 = env.addSource(consumer) val stream2 = stream1.map { x => ProductModel.readElement(Json.parse(x)) } stream2.print() env.execute("FlinkKafkaScala")
方案2:将解析逻辑放在case class伴生对象中
简化结构,把JSON解析逻辑直接放在case class的伴生对象里:
case class ProductJs(brand: String, category: String, productid: String, productname: String) object ProductJs { def fromJson(jsonStr: String): ProductJs = { val jsonElement = Json.parse(jsonStr) ProductJs( (jsonElement \ "brand").as[String], (jsonElement \ "category").as[String], (jsonElement \ "productid").as[String], (jsonElement \ "productname").as[String] ) } } val env = StreamExecutionEnvironment.getExecutionEnvironment val properties = new Properties() properties.setProperty("bootstrap.servers", "localhost:9092") properties.setProperty("group.id", "flink-consumer-group") properties.setProperty("auto.offset.reset", "earliest") val consumer = new FlinkKafkaConsumer[String]("my-topic", new SimpleStringSchema(), properties) val stream1 = env.addSource(consumer) val stream2 = stream1.map(ProductJs.fromJson) stream2.print() env.execute("FlinkKafkaScala")
内容的提问来源于stack exchange,提问作者Yaser
相关产品推荐
相关产品推荐

