Pulsar 2.7.1 Function发送AVRO消息未更新Schema导致反序列化失败
问题原因说明
两种发送逻辑的核心差异
你遇到的问题是Pulsar 2.7.x版本的Function框架特有设计导致的,两种发送逻辑的差异如下:
- 普通Spring Boot生产者:调用
newProducer创建生产者时,客户端会先把传入的Avro Schema和Broker端主题现有Schema做兼容性校验,开启set-is-allow-auto-update-schema的前提下,只要兼容性策略允许,会先注册新的Avro Schema到Broker,再完成生产者初始化,发送消息时自动关联最新的Schema版本。 - Function上下文发送:2.7.x版本的Function实例启动时,会预创建所有输出主题的内部生产者并缓存,预创建时默认读取主题当前已有的Schema(也就是你预创建主题时生成的JSON Schema)初始化生产者。后续调用
context.newOutputMessage时只会复用已缓存的生产者,你传入的AvroSchema.of(OurClass.class)参数实际不会生效,也不会触发新Schema的注册流程。
Schema不更新的根因
底层的消息发送逻辑确实复用了通用Producer的代码,但Function框架在之上做了生产者预创建的封装,跳过了普通生产者初始化阶段的Schema校验、注册逻辑,所以就算你传入了Avro Schema,也不会触发Broker端的Schema类型更新,主题始终保留最初生成的JSON Schema。
消费时TroubleFunction会按主题绑定的JSON Schema去反序列化Avro序列化后的二进制数据,自然会抛出JSON解析异常。
解决方案
按落地优先级推荐以下方案:
- 预创建主题时直接注册正确Schema
提前用pulsar admin CLI上传Avro Schema,从根源避免自动生成JSON Schema的问题,命令示例:
pulsar-admin schemas upload <tenant>/<namespace>/trouble-topic -f ./our-class-schema.json
其中our-class-schema.json需要明确指定type为AVRO,schema字段填写OurClass对应的Avro Schema定义。
- 部署Function时显式指定输出Schema
部署ValidationFunction时新增输出Schema配置参数,强制Function初始化生产者时用Avro Schema,触发Broker端Schema更新,参数示例:
pulsar-admin functions create \ --name ValidationFunction \ # 其他原有配置 --output-topic trouble-topic \ --output-schema-type AVRO \ --output-schema-class com.yourpackage.OurClass
- Function内部主动触发Schema注册
如果不想修改部署流程,可以在Function的open初始化方法中,手动调用客户端API创建一次trouble-topic的Avro生产者,创建过程会自动完成Schema注册,后续再用context.newOutputMessage发送也不会有问题。
另外如果开启了Schema兼容性校验,需要确认当前策略允许JSON到Avro的不兼容升级,否则自动更新Schema会被Broker拒绝,必要时可以临时将主题的兼容性策略调整为DISABLE,Schema更新完成后再改回原有策略。
该问题在Pulsar 2.8及以上版本已经修复,Function的newOutputMessage会正确识别传入的Schema参数并触发更新,有条件的话也可以考虑升级版本。
内容的提问来源于stack exchange,提问作者ante_f
相关产品推荐
相关产品推荐

