You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Scala中用Apache Flink读取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的算子任务需要可序列化,才能分发到集群节点执行。此处问题根源:

  • productJs case class定义在ProductModel object内部,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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.10 15:19:55