【问题标题】:How to use Avro serialization for scala case classes with Flink 1.7?如何将 Avro 序列化用于带有 Flink 1.7 的 scala 案例类?
【发布时间】:2019-03-20 17:06:30
【问题描述】:

我们有一个用 Scala 编写的 Flink 作业,它使用案例类(由 avrohugger 从 avsc 文件生成)来表示我们的状态。我们想使用 Avro 来序列化我们的状态,以便在我们更新模型时进行状态迁移。我们了解到,因为 Flink 1.7 Avro 序列化支持 OOTB。我们将 flink-avro 模块添加到类路径中,但是当从保存的快照恢复时,我们注意到它仍在尝试使用 Kryo 序列化。相关代码sn-p

case class Foo(id: String, timestamp: java.time.Instant)

val env = StreamExecutionEnvironment.getExecutionEnvironment
val conf = env.getConfig
conf.disableForceKryo()
conf.enableForceAvro()

val rawDataStream: DataStream[String] = env.addSource(MyFlinkKafkaConsumer)

val parsedDataSteam: DataStream[Foo] = rawDataStream.flatMap(new JsonParser[Foo])

// do something useful with it

env.execute("my-job")

在Foo 上执行状态迁移时(例如,通过添加字段并部署作业)我看到它尝试使用 Kryo 反序列化,这显然失败了。如何确保正在使用 Avro 序列化?

更新

发现了 https://issues.apache.org/jira/browse/FLINK-10897,因此仅从 1.8 afaik 开始支持使用 Avro 进行 POJO 状态序列化。我尝试使用最新的 1.8 RC 和一个从 SpecificRecord 扩展的简单 WordCount POJO:

/** MACHINE-GENERATED FROM AVRO SCHEMA. DO NOT EDIT DIRECTLY */
import scala.annotation.switch

case class WordWithCount(var word: String, var count: Long) extends 
  org.apache.avro.specific.SpecificRecordBase {
  def this() = this("", 0L)
  def get(field$: Int): AnyRef = {
    (field$: @switch) match {
      case 0 => {
        word
      }.asInstanceOf[AnyRef]
      case 1 => {
        count
      }.asInstanceOf[AnyRef]
      case _ => new org.apache.avro.AvroRuntimeException("Bad index")
    }
  }
  def put(field$: Int, value: Any): Unit = {
    (field$: @switch) match {
      case 0 => this.word = {
        value.toString
      }.asInstanceOf[String]
      case 1 => this.count = {
        value
      }.asInstanceOf[Long]
      case _ => new org.apache.avro.AvroRuntimeException("Bad index")
    }
    ()
  }
  def getSchema: org.apache.avro.Schema = WordWithCount.SCHEMA$
}

object WordWithCount {
     val SCHEMA$ = new org.apache.avro.Schema.Parser().parse(" . 
       {\"type\":\"record\",\"name\":\"WordWithCount\",\"fields\": 
       [{\"name\":\"word\",\"type\":\"string\"}, 
       {\"name\":\"count\",\"type\":\"long\"}]}")
}

然而,这也不是开箱即用的。然后我们尝试使用 flink-avro 的 AvroTypeInfo 定义我们自己的类型信息,但这失败了,因为 Avro 在类中查找 SCHEMA$ 属性 (SpecificData:285) 并且无法使用 Java 反射来识别 Scala 伴生对象中的 SCHEMA$ .

【问题讨论】:

    标签: scala apache-flink avro flink-streaming


    【解决方案1】:

    I could never get reflection to work 因为 Scala 的字段在后台是私有的。 AFAIK 唯一的解决方案是更新 Flink 以在 AvroInputFormat (compare) 中使用 avro 的非反射构造函数。

    在紧要关头,除了 Java,可以回退到 avro 的 GenericRecord,也许使用 avro4s 从 avrohugger 的 Standard 格式生成它们(注意 Avro4s 会从生成的 Scala 类型生成它自己的模式)

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-07-02
      • 1970-01-01
      • 2017-03-29
      • 1970-01-01
      • 2019-03-06
      • 1970-01-01
      • 2019-08-03
      相关资源
      最近更新 更多