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

Confluent Go Schema Registry验证不生效问题排查

问题

目前使用Confluent Schema Registry对生产的数据做schema验证,但发送与已注册schema完全不符的数据时,仍能成功发送,使用Go语言开发。

已注册的schema:

{
  "properties": {
    "device_id": {
      "type": "string"
    },
    "event_name": {
      "type": "string"
    },
    "platform": {
      "enum": [
        "android",
        "ios",
        "web"
      ],
      "type": "string"
    },
    "user_id": {
      "type": "string"
    }
  },
  "required": [
    "event_name",
    "user_id",
    "device_id",
    "platform"
  ],
  "title": "view_home",
  "type": "object"
}

实际发送的数据:

{
  "name": "First user",
  "favorite_number": "42",
  "favorite_color": "blue"
}

Go生产数据的代码:

func ProduceData(topic string, data []byte) {

    var d interface{}

    switch topic {
    case "view_home":
        d = ViewHome{}
        json.Unmarshal(data, &d)
    case "view_searchResult":
        d = ViewSearchResult{}
        json.Unmarshal(data, &d)
    }

    conf := ReadConfig()

    p, err := kafka.NewProducer(&conf)

    if err != nil {
        fmt.Printf("Failed to create producer: %s", err)
        os.Exit(1)
    }

    defer p.Close()

    if err != nil {
        fmt.Printf("Failed to create producer: %s\n", err)
        os.Exit(1)
    }

    fmt.Printf("Created Producer %v\n", p)

    client, err := schemaregistry.NewClient(schemaregistry.NewConfig("http://localhost:8081"))

    if err != nil {
        fmt.Printf("Failed to create schema registry client: %s\n", err)
        os.Exit(1)
    }

    serdeConfig := jsonschema.NewSerializerConfig()
    serdeConfig.AutoRegisterSchemas = false
    serdeConfig.UseLatestVersion = true
    serdeConfig.EnableValidation = true

    ser, err := jsonschema.NewSerializer(client, serde.ValueSerde, serdeConfig)

    if err != nil {
        fmt.Printf("Failed to create serializer: %s\n", err)
        os.Exit(1)
    }

    // Optional delivery channel, if not specified the Producer object's
    // .Events channel is used.
    deliveryChan := make(chan kafka.Event)
    value := User{
        Name:           "First user",
        FavoriteNumber: "42",
        FavoriteColor:  "blue",
    }
    payload, err := ser.Serialize(topic, &value)
    if err != nil {
        fmt.Printf("Failed to serialize payload: %s\n", err)
        os.Exit(1)
    }
    err = p.Produce(&kafka.Message{
        TopicPartition: kafka.TopicPartition{Topic: &topic, Partition: kafka.PartitionAny},
        Value:          payload,
    }, deliveryChan)
    if err != nil {
        fmt.Printf("Produce failed: %v\n", err)
        os.Exit(1)
    }

    e := <-deliveryChan
    m := e.(*kafka.Message)

    if m.TopicPartition.Error != nil {
        fmt.Printf("Delivery failed: %v\n", m.TopicPartition.Error)
    } else {
        fmt.Printf("Delivered message to topic %s [%d] at offset %v\n",
            *m.TopicPartition.Topic, m.TopicPartition.Partition, m.TopicPartition.Offset)
    }

    close(deliveryChan)
}

序列化配置:

serdeConfig := jsonschema.NewSerializerConfig()
serdeConfig.AutoRegisterSchemas = false
serdeConfig.UseLatestVersion = true
serdeConfig.EnableValidation = true

疑问:为何schema不匹配仍能正常序列化并发送?


分析与解决

核心问题点

  1. 代码逻辑错误:未使用Topic对应结构体序列化
    代码开头虽根据topic初始化了ViewHome等对应结构并反序列化数据到d,但后续完全未使用d,而是硬编码创建User实例进行序列化。这导致实际发送的数据与Topic绑定的schema完全无关,是逻辑层面的根本性错误。

  2. Schema验证未生效的潜在原因

    • 验证模式未开启严格检查:Confluent的JSON Schema序列化器默认可能采用宽松模式,仅做schema兼容性检查而非实例级验证。若缺少必填字段、字段类型不匹配等情况未触发错误,需确认是否需显式开启严格验证。
    • Subject名称不匹配:序列化器默认使用{topic}-value作为schema的subject,若Registry中注册的schema使用了自定义subject名称,会导致序列化器找不到对应schema,跳过验证。
    • SDK验证逻辑误解:EnableValidation=true可能仅开启schema兼容性检查(即检查生成的结构体schema与Registry中schema是否兼容),而非验证实际数据实例是否符合schema。若你的User结构体生成的schema被误判为兼容,就会跳过实例验证。

解决步骤

  1. 修正代码逻辑,使用Topic对应结构体
    移除硬编码的User实例,改用根据topic初始化的d变量进行序列化,确保数据结构与Topic绑定的schema一致:

    // 替换原来的value := User{...}代码块
    payload, err := ser.Serialize(topic, d)
    if err != nil {
        fmt.Printf("Failed to serialize payload: %s\n", err)
        os.Exit(1)
    }
    

    同时确保ViewHome结构体的JSON标签与schema字段完全匹配:

    type ViewHome struct {
        DeviceID  string `json:"device_id"`
        EventName string `json:"event_name"`
        Platform  string `json:"platform"`
        UserID    string `json:"user_id"`
    }
    
  2. 配置严格验证模式
    若使用的SDK支持,开启严格验证确保实例数据完全符合schema:

    serdeConfig := jsonschema.NewSerializerConfig()
    serdeConfig.AutoRegisterSchemas = false
    serdeConfig.UseLatestVersion = true
    serdeConfig.EnableValidation = true
    serdeConfig.ValidationStrict = true // 开启严格实例验证(需SDK支持)
    
  3. 确认Subject与Topic的对应关系
    检查Registry中schema的subject是否为view_home-value(ValueSerde默认格式),若使用自定义subject,需指定名称策略:

    serdeConfig.SubjectNameStrategy = schemaregistry.TopicNameStrategy
    
  4. 验证实例级检查是否生效
    若SDK默认仅做schema兼容性检查,可手动对序列化后的JSON进行实例验证,或升级SDK版本以支持完整的实例级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.08 15:15:55