Kafka Schema Registry验证异常:Broker校验记录失败求助
问题诊断与解决
核心原因
直接发送原始JSON数据,而开启Broker端Schema验证后,要求消息必须是经过Schema Registry序列化的格式——包含Schema ID的封装消息,而非纯JSON字符串。Broker无法识别原始JSON结构,因此抛出Broker: Broker failed to verify record.错误。
解决步骤
1. 补充Producer的Schema Registry配置
在ReadConfig()函数中添加Schema Registry地址及序列化器配置:
func ReadConfig() kafka.ConfigMap { return kafka.ConfigMap{ "bootstrap.servers": "your-kafka-broker:9092", // 替换为你的Kafka地址 "schema.registry.url": "your-schema-registry:8081", // 必须添加,指向你的Schema Registry "value.serializer": kafka.NewJSONSchemaSerializer, // 指定JSON Schema序列化器 "value.subject.name.strategy": "io.confluent.kafka.serializers.subject.TopicNameStrategy", // 按主题生成Subject,可根据需求调整 // 保留你原有的其他配置 } }
2. 修改消息发送逻辑
不再直接传递原始JSON字节,改为传递结构化数据,由序列化器自动完成Schema校验、注册(若未存在)及消息封装:
// 调整ProduceData参数,接受任意结构化数据而非[]byte func ProduceData(topic string, data interface{}) { conf := ReadConfig() p, err := kafka.NewProducer(&conf) if err != nil { fmt.Printf("Failed to create producer: %s", err) os.Exit(1) } defer p.Close() // 原有事件处理协程逻辑保持不变 go func() { for e := range p.Events() { switch ev := e.(type) { case *kafka.Message: if ev.TopicPartition.Error != nil { fmt.Printf("Failed to deliver message: %v\n", ev.TopicPartition) } else { fmt.Printf("Produced event to topic %s: key = %-10s value = %s\n", *ev.TopicPartition.Topic, string(ev.Key), string(ev.Value)) } } } }() // 直接传入结构化数据,序列化器会自动处理 p.Produce(&kafka.Message{ TopicPartition: kafka.TopicPartition{Topic: &topic, Partition: kafka.PartitionAny}, Value: data, }, nil) p.Flush(15 * 1000) } // main函数中传递结构体数据 func main() { // 定义与Schema匹配的结构体 type TestPayload struct { Test string `json:"test"` } kafka.ProduceData("schema_test", TestPayload{Test: "test1"}) }
3. 确认Schema已注册
确保你提供的JSON Schema已成功注册到Schema Registry,且Subject名称与Producer配置的策略一致(默认按{topic}-value生成,即此处为schema_test-value)。
关键提示
- Broker端的Schema验证依赖消息中携带的Schema ID,序列化器会自动与Registry交互,获取或注册Schema后将ID与数据封装发送。
- 若必须发送原始JSON,需关闭Broker端的
confluent.value.schema.validation配置,但这会失去Schema校验的防护能力。
内容的提问来源于stack exchange,提问作者gwkim
相关产品推荐
相关产品推荐

