【发布时间】: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