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不匹配仍能正常序列化并发送?
分析与解决
核心问题点
代码逻辑错误:未使用Topic对应结构体序列化
代码开头虽根据topic初始化了ViewHome等对应结构并反序列化数据到d,但后续完全未使用d,而是硬编码创建User实例进行序列化。这导致实际发送的数据与Topic绑定的schema完全无关,是逻辑层面的根本性错误。Schema验证未生效的潜在原因
- 验证模式未开启严格检查:Confluent的JSON Schema序列化器默认可能采用宽松模式,仅做schema兼容性检查而非实例级验证。若缺少必填字段、字段类型不匹配等情况未触发错误,需确认是否需显式开启严格验证。
- Subject名称不匹配:序列化器默认使用
{topic}-value作为schema的subject,若Registry中注册的schema使用了自定义subject名称,会导致序列化器找不到对应schema,跳过验证。 - SDK验证逻辑误解:
EnableValidation=true可能仅开启schema兼容性检查(即检查生成的结构体schema与Registry中schema是否兼容),而非验证实际数据实例是否符合schema。若你的User结构体生成的schema被误判为兼容,就会跳过实例验证。
解决步骤
修正代码逻辑,使用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"` }配置严格验证模式
若使用的SDK支持,开启严格验证确保实例数据完全符合schema:serdeConfig := jsonschema.NewSerializerConfig() serdeConfig.AutoRegisterSchemas = false serdeConfig.UseLatestVersion = true serdeConfig.EnableValidation = true serdeConfig.ValidationStrict = true // 开启严格实例验证(需SDK支持)确认Subject与Topic的对应关系
检查Registry中schema的subject是否为view_home-value(ValueSerde默认格式),若使用自定义subject,需指定名称策略:serdeConfig.SubjectNameStrategy = schemaregistry.TopicNameStrategy验证实例级检查是否生效
若SDK默认仅做schema兼容性检查,可手动对序列化后的JSON进行实例验证,或升级SDK版本以支持完整的实例级schema校验。
内容的提问来源于stack exchange,提问作者gwkim
相关产品推荐
相关产品推荐

