向带Schema的Kafka Topic发消息时遇Schema兼容错误排查
Kafka JSON Schema消息发送失败排查与解决
初始代码
String topic = "gcp.test_schema"; final Properties props = new Properties(); InputStream inputStream = Files.newInputStream(Paths.get("src/main/resources/client.properties")); props.load(inputStream); // 创建Kafka生产者 Producer<String, Object> producer = new KafkaProducer<>(props); String jsonPayload = "{\"age\": \"30\"}"; // 创建包含JSON负载的Kafka记录 ProducerRecord<String, Object> kafkaRecord = new ProducerRecord<>(topic, jsonPayload); // 发送记录到Kafka Topic producer.send(kafkaRecord, (metadata, exception) -> { if (exception != null) { exception.printStackTrace(); } else { System.out.println("消息发送成功!"); System.out.println("分区: " + metadata.partition()); System.out.println("偏移量: " + metadata.offset()); } });
Topic详情
❯ ./bin/confluent kafka topic describe gcp.test_schema Name | Value | Read-Only ------------------------------------------+----------------------------------------------------------+------------ cleanup.policy | delete | false compression.type | producer | true confluent.key.schema.validation | false | false confluent.key.subject.name.strategy | io.confluent.kafka.serializers.subject.TopicNameStrategy | false confluent.value.schema.validation | false | false confluent.value.subject.name.strategy | io.confluent.kafka.serializers.subject.TopicNameStrategy | false
已注册的Topic Schema
{ "properties": { "age": { "type": "string" } } }
生产者配置文件
# Kafka生产者、消费者和管理员所需的连接配置 bootstrap.servers=xxx.northamerica-northeast1.gcp.confluent.cloud:9092 security.protocol=SASL_SSL sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username='xxx' password='xxx/iQEXntX/xx'; sasl.mechanism=PLAIN key.serializer=org.apache.kafka.common.serialization.StringSerializer value.serializer=io.confluent.kafka.serializers.json.KafkaJsonSchemaSerializer confluent.value.schema.validation=true # Confluent Cloud Schema Registry所需的连接配置 schema.registry.url=https://qqqq.us-east1.gcp.confluent.cloud basic.auth.credentials.source=USER_INFO basic.auth.user.info=ffff:ssd+ddd+YySi
初始错误信息
org.apache.kafka.common.errors.InvalidConfigurationException: Schema being registered is incompatible with an earlier schema for subject "gcp.test_schema-value", details: [{errorType:"TYPE_CHANGED", description:"A type at path '#/' is different between the new schema and the old schema"}, {oldSchemaVersion: 1}, {oldSchema: '{"properties":{"age":{"type":"string"}}}'}, {compatibility: 'BACKWARD'}]; error code: 409; error code: 409
尝试的解决方法及对应错误
1. 使用Map发送
Producer<String, Map> producer = new KafkaProducer<String, Map>(props); ProducerRecord<String, Map> kafkaRecord = new ProducerRecord<>(topic, Map.of("age", "30"));
错误信息:
org.apache.kafka.common.errors.InvalidConfigurationException: Schema being registered is incompatible with an earlier schema for subject "gcp.test_schema-value", details: [{errorType:"ADDITIONAL_PROPERTIES_NARROWED", description:"An array or combined type at path '#/additionalProperties' has fewer elements in the new schema than the old schema"}, {oldSchemaVersion: 1}, {oldSchema: '{"properties":{"age":{"type":"string"}}}'}, {compatibility: 'BACKWARD'}]; error code: 409
2. 使用自定义类发送
class User { @JsonProperty public String age; public User() {} public User(String age) { this(age, null); } public User(String age, Object o) { } } Producer<String, User> producer = new KafkaProducer<String, User>(props); User user = new User("30"); ProducerRecord<String, User> kafkaRecord = new ProducerRecord<>(topic, user);
错误信息:
org.apache.kafka.common.errors.InvalidConfigurationException: Schema being registered is incompatible with an earlier schema for subject "gcp.test_schema-value", details: [{errorType:"ADDITIONAL_PROPERTIES_REMOVED", description:"The keyword at path '#/additionalProperties' in the new schema is not present in the old schema"}, {oldSchemaVersion: 1}, {oldSchema: '{"properties":{"age":{"type":"string"}}}'}, {compatibility: 'BACKWARD'}]; error code: 409; error code: 409
3. 修改Schema后的情况
将已注册的Schema修改为:
{ "$id": "http://example.com/myURI.schema.json", "$schema": "http://json-schema.org/draft-07/schema#", "additionalProperties": false, "description": "Sample schema to help you get started.", "properties": { "age": { "description": "The string type is used for strings of text.", "type": "string" } }, "title": "SampleRecord", "type": "object" }
此时发送无报错,但使用Map或自定义类时会自动生成新的Schema版本,消息未送达Topic也无错误提示。Map发送后生成的Schema版本如下:
{ "$schema": "http://json-schema.org/draft-07/schema#", "additionalProperties": {}, "title": "Map 1", "type": "object" }
问题根源
- 初始代码类型不匹配:直接发送JSON字符串时,
KafkaJsonSchemaSerializer会将其识别为string类型,但已注册的Schema未指定type: object,默认允许任意类型。由于Schema Registry启用了BACKWARD兼容性检查,新的string类型Schema无法兼容旧的“任意类型”Schema,导致TYPE_CHANGED错误。 - 自动生成Schema的兼容性问题:使用Map或自定义类时,序列化器会自动生成包含
additionalProperties配置的Schema,与已注册的无additionalProperties约束的Schema不兼容,触发BACKWARD兼容性检查失败。 - 消息未送达原因:修改Schema后,虽然无报错,但自动生成的新Schema与当前使用的Schema不匹配,且未调用
flush()和close()确保消息提交,导致消息未写入Kafka。
正确解决步骤
1. 显式指定要使用的Schema
在生产者配置中直接指定已注册的Schema,避免自动生成不兼容的版本:
import io.confluent.kafka.serializers.json.JsonSchema; import io.confluent.kafka.serializers.json.KafkaJsonSchemaSerializerConfig; // ... 加载配置代码 ... // 显式指定已注册的Schema String schemaStr = "{\"properties\":{\"age\":{\"type\":\"string\"}},\"type\":\"object\"}"; JsonSchema schema = JsonSchema.parse(schemaStr); props.put(KafkaJsonSchemaSerializerConfig.VALUE_SCHEMA_CONFIG, schema); // 使用自定义类发送 Producer<String, User> producer = new KafkaProducer<>(props); User user = new User("30"); ProducerRecord<String, User> kafkaRecord = new ProducerRecord<>(topic, user); // 发送后务必刷新并关闭生产者 producer.send(kafkaRecord, (metadata, exception) -> { if (exception != null) { exception.printStackTrace(); } else { System.out.println("消息发送成功!"); System.out.println("分区: " + metadata.partition()); System.out.println("偏移量: " + metadata.offset()); } }); producer.flush(); producer.close();
2. 调整自定义类的Schema生成
确保自定义类生成的Schema与已注册版本一致,可通过注解控制:
import com.fasterxml.jackson.annotation.JsonIgnoreProperties; @JsonIgnoreProperties(ignoreUnknown = false) // 禁止未知属性,匹配Schema的additionalProperties: false class User { public String age; public User() {} public User(String age) { this.age = age; } }
3. 修改Schema兼容性级别(可选,不推荐生产环境)
如果不需要严格的BACKWARD兼容性,可将Schema Registry的兼容性级别改为NONE:
./bin/confluent schema-registry config set --compatibility NONE --subject gcp.test_schema-value
4. 确保生产者正确提交消息
发送消息后必须调用flush()和close(),确保消息被提交到Kafka集群。
内容的提问来源于stack exchange,提问作者Gaurang Shah
相关产品推荐
相关产品推荐

