You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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反序列化器

  1. 引入Confluent Schema Registry依赖(版本需与生产者保持一致):
<!-- Maven示例 -->
<dependency>
    <groupId>io.confluent</groupId>
    <artifactId>kafka-protobuf-serializer</artifactId>
    <version>${confluent.version}</version>
</dependency>
  1. 修改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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 18:23:19