Spark使用V2 Schema读取Avro V1数据时向后兼容性失效问题
问题根因
Avro的向后兼容能力生效的前提是解码器同时感知写入端Schema(Writer Schema,即写数据用的V1)和读取端Schema(Reader Schema,即读数据用的V2),只有这样Avro才会自动将旧Schema写入的数据适配到新Schema结构,为缺失的新增字段填充默认值。
你当前的调用方式直接给from_avro传入了单个Schema,Spark默认会将这个Schema同时作为Writer Schema和Reader Schema使用:如果你传的是V2,相当于告诉Avro解码器待解析的二进制是用V2写入的,需要读取3个字段的内容,但实际二进制是V1写入的,只有2个字段的字节长度,解码器读到末尾找不到第三个字段的数据就抛出EOF异常。
解决方法
- 分别加载写入端V1 Schema和读取端V2 Schema的内容
- 调用
from_avro时显式指定Reader Schema参数,示例代码如下:
import org.apache.spark.sql.avro.functions._ import java.nio.file.{Paths, Files} // 加载写入端用的V1 Schema val writerSchema = new String(Files.readAllBytes(Paths.get("./examples/src/main/resources/V1.avsc"))) // 加载读取端用的V2 Schema val readerSchema = new String(Files.readAllBytes(Paths.get("./examples/src/main/resources/V2.avsc"))) val df = spark .readStream .format("kafka") .option("kafka.bootstrap.servers", "host1:port1,host2:port2") .option("subscribe", "topic1") .load() // 调用from_avro时通过options指定readerSchema val output = df .select(from_avro($"value", writerSchema, Map("readerSchema" -> readerSchema)) as $"avroFields")
补充说明
- 你提供的V1、V2 Schema本身满足向后兼容要求:新增字段
three带null默认值,原有字段顺序、名称、类型、Schema的name和namespace都未修改,配置正确后可以正常读取V1写入的数据,three字段会自动填充为null。 - 如果写入端后续也会升级到V2 Schema写入数据,不需要修改读取端代码:只要写入端用的Schema和你传的writerSchema兼容,都可以正常解码适配到readerSchema的结构。
内容的提问来源于stack exchange,提问作者darkknight444
相关产品推荐
相关产品推荐

