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

Kafka Schema Registry如何拦截不符合已注册Schema规范的消息

Kafka生产端Avro Schema强校验拦截实现方案

你要的消息结构不匹配就直接拦截、禁止写入的效果,不需要自定义拦截器,直接用Confluent官方Avro序列化组件配合正确配置就能原生实现,校验逻辑在客户端序列化阶段执行,不符合规则的消息根本不会发送到Broker。

前置准备

必须使用Confluent提供的KafkaAvroSerializer作为消息值的序列化器,不要自己手动把对象转成字节数组发送,否则序列化器无法读取消息结构做校验。
以Java客户端为例,先引入和你集群版本匹配的依赖:

<dependency>
    <groupId>io.confluent</groupId>
    <artifactId>kafka-avro-serializer</artifactId>
    <version>对应Confluent平台版本号</version>
</dependency>

其他语言客户端(Python/Go/.NET)只要用官方维护的Confluent Kafka客户端,都自带这套序列化校验能力,配置逻辑通用。

核心配置(少一个都达不到强校验效果)

生产者端添加以下参数:

  • value.serializer: io.confluent.kafka.serializers.KafkaAvroSerializer
  • schema.registry.url: http://localhost:38081 替换为你实际的Schema Registry地址
  • auto.register.schemas: false 关键配置:禁止生产者自动注册新Schema,只会拉取你提前上传到Registry的Schema做匹配
  • use.latest.version: true 校验时拉取对应subject下最新的已注册Schema作为校验基准
  • 注意subject命名匹配:你当前注册Schema的subject名是teste,默认序列化器的命名规则是[Topic名]-value,如果你的目标Topic名是teste,建议把已注册的subject名改为teste-value适配默认规则,避免拉不到Schema报错。

校验生效逻辑

配置完成后,发送消息有两种方式,都会自动触发强校验:

  1. 用Avro工具根据你定义的Schema生成对应data.brado.operacao.Operacao类,发送时直接传该类的实例,所有字段按类型赋值即可
  2. 传通用GenericRecord对象,手动按Schema定义填字段名和字段值

只要出现以下任意一种情况,序列化阶段会直接抛出SerializationException,终止发送流程,消息不会写入Broker:

  • 缺失任意一个定义的字段(你当前Schema里9个字段全是无默认值的string类型,不允许缺省)
  • 字段类型不匹配(比如给id_operacao传数字而非字符串)
  • 携带Schema里未定义的额外字段

正确发送示例(Java)

Properties props = new Properties();
props.put("bootstrap.servers", "你的Kafka Broker地址");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "io.confluent.kafka.serializers.KafkaAvroSerializer");
props.put("schema.registry.url", "http://localhost:38081");
props.put("auto.register.schemas", "false");
props.put("use.latest.version", "true");

KafkaProducer<String, Operacao> producer = new KafkaProducer<>(props);

// 合规消息:9个字段全为string类型,无多余字段
Operacao validMsg = Operacao.newBuilder()
    .setIdOperacao("OP001")
    .setTipoContainer("40HQ")
    .setDescricaoOperacao("重箱出场")
    .setEntrega("1")
    .setColeta("0")
    .setDescricaoChecklist("核对箱号、封条号完好")
    .setCheio("1")
    .setAtivo("1")
    .setTipoOperacao("OUT")
    .build();
// 该消息校验通过,正常发送
producer.send(new ProducerRecord<>("teste", validMsg));

// 违规消息示例:比如漏传ativo字段、给cheio传布尔值true而非字符串"1"
// 序列化阶段直接抛错,消息不会到达Broker

内容的提问来源于stack exchange,提问作者souzatorquato

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 22:12:18