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

Kafka Connect报错:Schema与postgres.public.customers-value早期版本不兼容

问题描述

我刚接触Kafka与ksqldb,正尝试用ksqldb搭建数据管道。已创建如下Postgres源Connector:

CREATE SOURCE CONNECTOR postgres_reader WITH (
    'connector.class' = 'io.debezium.connector.postgresql.PostgresConnector',
    'database.hostname' = 'host.docker.internal',
    'database.port' = '5432',
    'database.user' = 'postgres-user',
    'database.password' = 'postgres-pw',
    'database.dbname' = 'customers',
    'database.server.name' = 'customers',
    'table.whitelist' = 'public.*',
     'value.converter' = 'io.confluent.connect.avro.AvroConverter',
    'key.converter' = 'io.confluent.connect.avro.AvroConverter',
    'value.converter.schema.registry.url' = 'http://host.docker.internal:8085',
    'key.converter.schema.registry.url' = 'http://host.docker.internal:8085',
    'transforms' = 'unwrap,createKey,ExtractField',
    'transforms.createKey.type' = 'org.apache.kafka.connect.transforms.ValueToKey',
    'transforms.createKey.fields' = 'id',
    'transforms.ExtractField.type' = 'org.apache.kafka.connect.transforms.ExtractField$Key',
    'transforms.ExtractField.field' = 'id',
    'transforms.unwrap.type' = 'io.debezium.transforms.ExtractNewRecordState',
    'transforms.unwrap.drop.tombstones' = 'false',
    'transforms.unwrap.delete.handling.mode' = 'rewrite',
    'topic.prefix' = 'postgres'
);

Postgres数据库中有customers表:

CREATE TABLE IF NOT EXISTS public.customers
(
    id text NOT NULL,
    name text,
    age integer,
    CONSTRAINT customers_pkey PRIMARY KEY (id)
)

同时在ksqldb中创建了流:

CREATE STREAM customers (
    id VARCHAR key,
    name VARCHAR,
    age INTEGER
) WITH (
    kafka_topic = 'postgres.public.customers',
    value_format = 'avro',
    PARTITIONS=1
);

创建流后自动生成了对应的Kafka主题及Schema Registry中的值与键Schema。但向customers表插入数据时,Connector触发如下Schema兼容性错误:

Caused by: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: Schema being registered is incompatible with an earlier schema for subject "postgres.public.customers-value", details: [{errorType:'NAME_MISMATCH', description:'The name of the schema has changed (path '/name')', additionalInfo:'expected: io.confluent.ksql.avro_schemas.KsqlDataSourceSchema'}, {errorType:'READER_FIELD_MISSING_DEFAULT_VALUE', description:'The field 'id' at path '/fields/0' in the new schema has no default value and is missing in the old schema', additionalInfo:'id'}, {oldSchemaVersion: 1}, {oldSchema: '{"type":"record","name":"KsqlDataSourceSchema","namespace":"io.confluent.ksql.avro_schemas","fields":[{"name":"NAME","type":["null","string"],"default":null},{"name":"AGE","type":["null","int"],"default":null}],"connect.name":"io.confluent.ksql.avro_schemas.KsqlDataSourceSchema"}'}, {compatibility: 'BACKWARD'}]; error code: 409
        at io.confluent.kafka.schemaregistry.client.rest.RestService.sendHttpRequest(RestService.java:297)
        at io.confluent.kafka.schemaregistry.client.rest.RestService.httpRequest(RestService.java:367)
        at io.confluent.kafka.schemaregistry.client.rest.RestService.registerSchema(RestService.java:544)
        at io.confluent.kafka.schemaregistry.client.rest.RestService.registerSchema(RestService.java:532)
        at io.confluent.kafka.schemaregistry.client.rest.RestService.registerSchema(RestService.java:490)

疑问:

  1. 为何Connector在向主题发送新消息时要注册新Schema?
  2. 若注册新Schema,其结构应与现有Schema一致,为何仍出现兼容性错误?

解答

问题1:Connector注册新Schema的原因

  • 你先通过ksqlDB创建了流,此时ksqlDB会自动在Schema Registry中为postgres.public.customers-value主题生成一个初始Schema(即错误信息里的KsqlDataSourceSchema)。
  • Debezium Connector捕获Postgres变更数据后,会根据数据库表结构生成对应的Avro Schema,这个Schema与ksqlDB预先创建的Schema完全不同,所以Connector需要尝试将新Schema注册到Registry,触发了兼容性检查。

问题2:结构看似一致却不兼容的核心原因

从错误信息可提取两个关键不兼容点:

  1. Schema全名不匹配:ksqlDB生成的Schema全名是io.confluent.ksql.avro_schemas.KsqlDataSourceSchema,而Debezium生成的Schema是基于数据库表的命名规则(比如包含postgres.public.customers相关标识)。Schema Registry的BACKWARD兼容性规则会校验Schema的全名(名称+命名空间),名称不一致直接判定不兼容。
  2. 字段差异:ksqlDB创建的流中id被定义为Key,因此Value Schema里不包含id字段;但Debezium经过unwrap转换后,Value中包含无默认值的id字段。在BACKWARD模式下,新增无默认值的字段会导致不兼容——旧Schema的消费者无法处理这个新增的必填字段。

修复建议

  • 推荐方案:不要提前用ksqlDB创建流,先让Debezium Connector捕获数据并生成Schema,再基于已有的Kafka主题自动推断Schema创建流,命令示例:
    CREATE STREAM customers WITH (
        kafka_topic = 'postgres.public.customers',
        value_format = 'avro'
    );
    
  • 备选方案:调整Connector的transform配置,移除Value中的id字段(因为已作为Key),同时自定义Schema命名规则匹配ksqlDB的格式,但这种方式配置繁琐,不推荐。

内容的提问来源于stack exchange,提问作者Noy Gafni

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 18:07:07