Spark-avro报BufferHolder size negative错误的原因及排查方法
问题根因
这个报错是Spark执行from_avro反序列化写入UnsafeRow时,计算出的缓冲区扩容值为负数,触发了BufferHolder的合法性校验。结合你的场景,基本是两个问题叠加导致:
- 最核心的错误:你用字符串函数
substring裁剪二进制Avro数据,偏移计算完全错了。Spark SQL的substring最初是为处理字符串设计的,作用于Kafka读出的二进制value列(类型为Array[Byte])时,会先按UTF-8编码把二进制转成字符串再做截断,多字节的Avro值会被错误识别、替换,最后得到的valueWithoutEmbeddedInfo根本不是合法的Avro二进制内容,反序列化时解析出的字段长度完全错乱,自然会算出负的缓冲区大小。 - 第二个诱因:你用的单份Schema大小超过5MB,刚好触发Spark 2.4版本
from_avro的已知缺陷。Spark 2.4的Avro反序列化器生成UnsafeRow时,会先根据Schema结构预估整行占的内存大小,当Schema嵌套深、字段总数太多时,预估值会超过Int类型上限溢出成负数,哪怕Avro数据本身合法,也会抛这个扩容错误。
你提到Kafka UI可以正常反序列化消息,说明Kafka里存的原始消息本身没问题,故障完全出在Spark侧的处理逻辑上。
排查&解决步骤
先修正头部裁剪逻辑,绝对不能用字符串函数处理二进制数据,改用二进制操作函数裁剪前13字节:
把你原来的substring逻辑替换成:import org.apache.spark.sql.functions.expr .withColumn("valueWithoutEmbeddedInfo", expr("slice(value, 14, length(value)-13)"))slice作用在二进制列/数组列时是严格按字节/元素位置裁剪,不会做编码转换,能保证拿到干净的Avro二进制内容。
验证方式很简单:拿1条消息,打印裁剪前后的字节长度,差值必须严格等于13,差值不对就说明裁剪逻辑有问题。解决大Schema触发的Spark 2.4兼容问题:
- 优先方案:别把5MB的全量Schema直接传给
from_avro。Avro本身支持Schema演进,你只需要把实际要用到的字段抽出来生成精简的解析Schema就行,不需要传全量,反序列化时会自动忽略Schema里没定义的多余字段,既能解决溢出问题,还能大幅提升反序列化性能。 - 如果必须用全量Schema,先把Spark任务依赖的spark-avro包升级到2.4线的最新小版本2.4.8,这个版本修了大Schema下预估行大小溢出的问题;如果没法升级依赖,提交任务时加参数
spark.sql.avro.rowPreferSize=268435456(对应256MB),调大Avro解析的初始行缓冲区大小,绕过预估值计算溢出的逻辑。 - 快速验证方法:手动取一条Kafka消息,裁掉前13字节存成本地avro文件,用spark-shell本地读这个文件,传入你拿到的
jsonFormatedSchema做反序列化,如果本地复现同样的错误,就能确认是大Schema触发的版本bug。
- 优先方案:别把5MB的全量Schema直接传给
额外校验项
确认你从SchemaRegistry拿到的jsonFormatedSchema是标准Avro Schema格式,没有多余转义字符、外层包装字段。Hortonworks SchemaRegistry返回的SchemaText默认是带命名空间的标准Avro Schema,本身不会有问题,但如果中途做过字符串拼接处理,就要额外校验。
避坑提醒:不要尝试把二进制value先转成字符串再裁剪。Avro是二进制序列化格式,转字符串过程中遇到UTF-8无法识别的字节会被替换成占位符,直接导致Avro内容损坏,所有二进制处理必须用Spark原生的数组/二进制操作函数。
内容的提问来源于stack exchange,提问作者napkin-pumpkin
相关产品推荐
相关产品推荐

