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

向带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"
}

问题根源

  1. 初始代码类型不匹配:直接发送JSON字符串时,KafkaJsonSchemaSerializer会将其识别为string类型,但已注册的Schema未指定type: object,默认允许任意类型。由于Schema Registry启用了BACKWARD兼容性检查,新的string类型Schema无法兼容旧的“任意类型”Schema,导致TYPE_CHANGED错误。
  2. 自动生成Schema的兼容性问题:使用Map或自定义类时,序列化器会自动生成包含additionalProperties配置的Schema,与已注册的无additionalProperties约束的Schema不兼容,触发BACKWARD兼容性检查失败。
  3. 消息未送达原因:修改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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 17:29:50