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)
疑问:
- 为何Connector在向主题发送新消息时要注册新Schema?
- 若注册新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:结构看似一致却不兼容的核心原因
从错误信息可提取两个关键不兼容点:
- Schema全名不匹配:ksqlDB生成的Schema全名是
io.confluent.ksql.avro_schemas.KsqlDataSourceSchema,而Debezium生成的Schema是基于数据库表的命名规则(比如包含postgres.public.customers相关标识)。Schema Registry的BACKWARD兼容性规则会校验Schema的全名(名称+命名空间),名称不一致直接判定不兼容。 - 字段差异: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
相关产品推荐
相关产品推荐

