Kafka扩容Topic分区后报Broker failed to validate record如何解决
问题解答
分区扩容操作完成后,不需要手动执行任何刷新Broker缓存的额外操作。Kafka Controller会自动将最新的分区元数据同步到集群内所有Broker,客户端也会通过内置的元数据刷新机制自动感知新分区,整个流程是集群自动完成的,无需人工介入。
你碰到的Broker failed to validate record报错和分区扩容操作没有直接关联,触发原因就是Broker端开启的Schema校验逻辑拦截了不符合要求的消息:
- 当Topic开启了Broker端Schema校验(对应配置项为
confluent.key.schema.validation、confluent.value.schema.validation,值为true即开启),Broker会在消息写入阶段直接校验消息是否携带合法的Schema标识、消息内容是否匹配Schema Registry中注册的对应规则,校验不通过就会直接返回该错误,不会等消费阶段再做校验。 - 你当前用kcat发送的是无格式的纯文本
test,没有携带Confluent序列化协议要求的Schema标识头(消息开头固定的5字节结构:1字节魔数+4字节Schema ID),Broker无法识别到合法的Schema关联信息,就会直接抛出校验失败错误。你可以尝试往扩容前就存在的老分区发送同样的纯文本消息,会触发完全相同的报错。
可以按以下步骤排查修复:
- 先确认报错根因:临时关闭对应Topic的Broker端值Schema校验,再执行同样的生产测试命令:
如果关闭校验后消息可以正常写入,即可确认报错是Schema校验拦截导致,和分区扩容流程无关。kafka-configs --alter --entity-type topics --entity-name <你的Topic名称> --add-config confluent.value.schema.validation=false - 如果业务要求必须保留Broker端Schema校验,就不能直接发送无Schema的纯文本测试消息:要么使用集成了Schema序列化逻辑的正式业务客户端生产消息,要么在使用kcat时指定Schema Registry地址、按照对应序列化格式(Avro/Protobuf/JSON Schema)构造合法消息。
- 如果你确认扩容前执行相同的kcat命令可以正常发送消息,需要检查Terraform中定义的Kafka Topic资源配置:部分版本的Confluent Terraform Provider对Schema校验配置有默认开启的逻辑,可能是你调整分区数执行apply时,没有显式指定校验开关的配置,导致Provider顺带把校验配置改成了开启状态。
内容的提问来源于stack exchange,提问作者Elisha Choo
相关产品推荐
相关产品推荐

