Spring Cloud从PubSub读取Avro消息遇空指针异常求助
嘿,我仔细看了你的代码和错误栈,马上找到了问题根源:你的消息转换器在消费消息时居然在尝试序列化对象(写Avro),而不是反序列化!
看错误栈里的ReflectDatumWriter.write调用——这明显是序列化操作,但你明明是在消费PubSub里的Avro二进制消息,应该做的是反序列化字节数组到Session对象才对。问题出在你用的AvroSchemaMessageConverter配置不对,它默认更适配Spring Cloud Stream的生产者场景,没有明确配置成消费时的反序列化模式。
下面给你两种靠谱的解决方案:
方案1:自定义Avro反序列化转换器(最直观可靠)
如果你的Session是通过Avro schema文件生成的Specific类,直接写一个针对性的转换器:
private MessageConverter sessionMessageConverter() { return new SimpleMessageConverter() { @Override public Object fromMessage(Message<?> message, Class<?> targetClass) { if (Session.class.isAssignableFrom(targetClass)) { byte[] avroPayload = (byte[]) message.getPayload(); try { // 用SpecificDatumReader处理Avro生成的类 DatumReader<Session> reader = new SpecificDatumReader<>(Session.getClassSchema()); Decoder decoder = DecoderFactory.get().binaryDecoder(avroPayload, null); return reader.read(null, decoder); } catch (IOException e) { throw new MessageConversionException("Failed to deserialize Avro Session message", e); } } return super.fromMessage(message, targetClass); } }; }
如果你的Session是普通POJO(靠Avro反射机制处理),把SpecificDatumReader换成ReflectDatumReader<Session>就行:
DatumReader<Session> reader = new ReflectDatumReader<>(Session.class);
方案2:修正AvroSchemaMessageConverter的配置
要是你想继续用AvroSchemaMessageConverter,得明确告诉它要处理二进制消息,并且目标转换类型是Session:
private MessageConverter sessionMessageConverter() { // 指定处理二进制类型的消息 payload AvroSchemaMessageConverter converter = new AvroSchemaMessageConverter(MediaType.APPLICATION_OCTET_STREAM); converter.setSchema(Session.getClassSchema()); // 设置要转换的目标类型 converter.setTargetType(TypeDescriptor.valueOf(Session.class)); return converter; }
额外提醒
- 确保Spring应用里的
Session类和Dataflow中使用的完全一致(Schema、类版本),不然反序列化肯定会出问题。 - 你当前用的Spring Cloud GCP(1.0.0.M3)和Avro(1.8.2)版本都比较老,建议升级到稳定版(比如Spring Cloud GCP 1.2.x+搭配Avro 1.9.x+),能避免不少兼容性坑。
内容的提问来源于stack exchange,提问作者Grzes
相关产品推荐
相关产品推荐

