Reactor Kafka中无法使用Schema Registry的问题排查
问题描述
我正尝试使用Reactor Kafka搭建Kafka消费者,对应的生产者已集成Kafka Schema Registry。
消费者配置代码如下:
@Value("${spring.kafka.schemaRegistryUrls}") private String schemaRegistryEnvVarValue; @Bean public ReceiverOptions<String, MyProto> kafkaReceiverOptionsFloor( KafkaProperties kafkaProperties) { final Map<String, Object> kafkaConsumerProperties = kafkaProperties.buildConsumerProperties(); for (Map.Entry<String, KafkaProperties.Consumer> entry : kafkaConsumerPropertiesConfig.getConsumer().entrySet()) { if (kafkaTopics.contains(entry.getKey())) { kafkaConsumerProperties.putAll(entry.getValue().buildProperties()); } } kafkaConsumerProperties.put("schema.registry.url", schemaRegistryEnvVarValue); final ReceiverOptions<String, MyProto> basicReceiverOptions = ReceiverOptions.<String, MyProto>create( kafkaConsumerProperties) .withValueDeserializer(new MyProtoDeserializer()) // disabling auto commit, since we are managing committing once // record is // processed .commitInterval(Duration.ZERO) .commitBatchSize(0); kafkaConsumerProperties.forEach((k, v) -> log.debug("k2 {} v2 {}", k, v)); return basicReceiverOptions .subscription(kafkaTopicsFloor) .addAssignListener(partitions -> log.debug("onPartitionsAssigned {}", partitions)) .addRevokeListener(partitions -> log.debug("onPartitionsRevoked {}", partitions)); } @Bean public ReactiveKafkaConsumerTemplate<String, MyProto> reactiveKafkaConsumerTemplate( ReceiverOptions<String, MyProto> kafkaReceiverOptions) { return new ReactiveKafkaConsumerTemplate<>(kafkaReceiverOptions); }
运行时遇到异常:Protocol message contained an invalid tag (zero)。未使用Schema Registry的单元测试能够正常解析数据,看起来Schema Registry未被正确使用,请问我哪里出错了?
反序列化器代码如下:
@Slf4j public class MyProtoDeserializer implements Deserializer<MyProto> { public MyProtoDeserializer() {} /** * Deserializes the data to my_proto from byte array. * * @param topic * @param data * @return */ @Override public MyProto deserialize(final String topic, final byte[] data) { if (data == null) { return null; } // TODO: Use schemaregistry and kpow try { return MyProto.getDefaultInstance() .getParserForType() .parseFrom(data); } catch (Exception ex) { log.debug("Exception in MyProto parse {}", ex.getMessage()); return MyProto.getDefaultInstance(); } } }
解决方案
问题核心是你的自定义MyProtoDeserializer没有处理Schema Registry的消息格式。当生产者通过Schema Registry发送消息时,字节数据会包含1字节魔术位+4字节Schema ID的前缀,不是纯Protobuf结构,直接解析自然会触发无效标签错误。
推荐方案:使用Confluent官方Protobuf反序列化器
- 引入Confluent Schema Registry依赖(版本需与生产者保持一致):
<!-- Maven示例 --> <dependency> <groupId>io.confluent</groupId> <artifactId>kafka-protobuf-serializer</artifactId> <version>${confluent.version}</version> </dependency>
- 修改ReceiverOptions配置,替换自定义反序列化器为官方实现:
@Bean public ReceiverOptions<String, MyProto> kafkaReceiverOptionsFloor( KafkaProperties kafkaProperties) { // 原有属性配置逻辑不变 final Map<String, Object> kafkaConsumerProperties = kafkaProperties.buildConsumerProperties(); // ... 其他属性赋值逻辑 ... kafkaConsumerProperties.put("schema.registry.url", schemaRegistryEnvVarValue); // 初始化官方Protobuf反序列化器 ProtobufDeserializer<MyProto> protobufDeserializer = new ProtobufDeserializer<>(MyProto.class); protobufDeserializer.configure(kafkaConsumerProperties, false); final ReceiverOptions<String, MyProto> basicReceiverOptions = ReceiverOptions.<String, MyProto>create(kafkaConsumerProperties) .withValueDeserializer(protobufDeserializer) .commitInterval(Duration.ZERO) .commitBatchSize(0); // ... 订阅、监听器配置逻辑不变 ... return basicReceiverOptions; }
备选方案:修改自定义反序列化器(不推荐,重复造轮子)
如果坚持使用自定义实现,需要手动解析Schema Registry的消息前缀:
@Override public MyProto deserialize(final String topic, final byte[] data) { if (data == null) { return null; } try { // 校验消息长度,Schema Registry格式至少5字节(1位魔术+4位Schema ID) if (data.length < 5) { throw new IllegalArgumentException("Invalid message length for Schema Registry format"); } // 校验魔术位,Protobuf序列化器的魔术位固定为0 byte magicByte = data[0]; if (magicByte != 0) { throw new IllegalArgumentException("Unsupported magic byte: " + magicByte); } // 解析Schema ID int schemaId = ByteBuffer.wrap(data, 1, 4).getInt(); // 从Schema Registry获取对应Schema(需引入Schema Registry客户端依赖) SchemaRegistryClient client = new CachedSchemaRegistryClient(schemaRegistryEnvVarValue, 100); client.getById(schemaId); // 跳过前缀,解析真正的Protobuf数据(从第5字节开始) return MyProto.getDefaultInstance().getParserForType().parseFrom(data, 5, data.length -5); } catch (Exception ex) { log.error("Failed to deserialize MyProto", ex); return MyProto.getDefaultInstance(); } }
额外检查项
- 确认
schema.registry.url配置正确,消费者网络能访问到Schema Registry服务 - 确认生产者使用的是Confluent官方
ProtobufSerializer,而非纯Protobuf序列化 - 单元测试用的是纯Protobuf数据,所以能正常解析,但生产环境消息带Schema Registry前缀,导致解析失败
内容的提问来源于stack exchange,提问作者Gajukorse
相关产品推荐
相关产品推荐

