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.KafkaAvroSerializerschema.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报错。
校验生效逻辑
配置完成后,发送消息有两种方式,都会自动触发强校验:
- 用Avro工具根据你定义的Schema生成对应
data.brado.operacao.Operacao类,发送时直接传该类的实例,所有字段按类型赋值即可 - 传通用
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
相关产品推荐
相关产品推荐

