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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 02:05:31