使用Schema Registry反序列化Topic创建ksqlDB表出错求助
问题:ksqlDB查询JSON Schema注册表的表时出现序列化错误
环境配置
--- version: '2' services: zookeeper: image: confluentinc/cp-zookeeper:7.5.1 hostname: zookeeper container_name: zookeeper ports: - "2181:2181" environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 broker: image: confluentinc/cp-kafka:7.5.1 hostname: broker container_name: broker depends_on: - zookeeper ports: - "9092:9092" environment: KAFKA_BROKER_ID: 1 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:9092 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 KAFKA_SCHEMA_REGISTRY_URL: 'http://schema-registry:8081' KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' schema-registry: image: confluentinc/cp-schema-registry:7.5.1 hostname: schema-registry container_name: schema-registry depends_on: - broker ports: - "8081:8081" environment: SCHEMA_REGISTRY_HOST_NAME: schema-registry SCHEMA_REGISTRY_KAFKASTORE_CONNECTION_URL: 'zookeeper:2181' SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: PLAINTEXT://broker:9092 SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081 SCHEMA_REGISTRY_DEBUG: 'true' ksqldb-server: image: confluentinc/cp-ksqldb-server:7.5.1 hostname: ksqldb-server container_name: ksqldb-server depends_on: - broker - schema-registry ports: - "8088:8088" environment: KSQL_LISTENERS: http://0.0.0.0:8088 KSQL_BOOTSTRAP_SERVERS: broker:9092 KSQL_KSQL_SCHEMA_REGISTRY_URL: 'http://schema-registry:8081' KSQL_KSQL_LOGGING_PROCESSING_STREAM_AUTO_CREATE: "true" KSQL_KSQL_LOGGING_PROCESSING_TOPIC_AUTO_CREATE: "true" ksqldb-cli: image: confluentinc/cp-ksqldb-cli:7.5.1 container_name: ksqldb-cli depends_on: - ksqldb-server entrypoint: /bin/sh tty: true
操作流程与问题现象
1. Schema注册
注册了两个JSON Schema,但手动指定了相同的$id:
- Key Schema:
{ "$id": "http://schema-registry:8081/schemas/ids/1", "$schema": "https://json-schema.org/draft/2020-12/schema", "title": "Location", "type": "object", "properties": { "profileId": { "type":"string" } } }
提交文件schema-key-registry.json:
{ "schemaType":"JSON", "schema":"{\"$id\":\"http://schema-registry:8081/schemas/ids/1\",\"$schema\":\"https://json-schema.org/draft/2020-12/schema\",\"title\":\"Location\",\"type\":\"object\",\"properties\":{\"profileId\":{\"type\":\"string\"}}}" }
- Value Schema:
{ "$id": "http://schema-registry:8081/schemas/ids/1", "$schema": "https://json-schema.org/draft/2020-12/schema", "title": "Location", "type": "object", "properties": { "profileId": { "type": "string", "description": "The id of the location." }, "latitude": { "type": "number", "minimum": -90, "maximum": 90, "description": "The location's latitude." }, "longitude": { "type": "number", "minimum": -180, "maximum": 180, "description": "The location's longitude." } } }
提交文件schema-value-registry.json:
{ "schemaType":"JSON", "schema":"{\"$id\":\"http://schema-registry:8081/schemas/ids/1\",\"$schema\":\"https://json-schema.org/draft/2020-12/schema\",\"title\":\"Location\",\"type\":\"object\",\"properties\":{\"profileId\":{\"type\":\"string\",\"description\":\"The id of the location.\"},\"latitude\":{\"type\":\"number\",\"minimum\":-90,\"maximum\":90,\"description\":\"The location's latitude.\"},\"longitude\":{\"type\":\"number\",\"minimum\":-180,\"maximum\":180,\"description\":\"The location's longitude.\"}}}" }
执行注册命令:
curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \ --data "$(cat schema-key-registry.json)" \ http://localhost:8081/subjects/locations-key/versions {"id":1} curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \ --data "$(cat schema-value-registry.json)" \ http://localhost:8081/subjects/locations/versions {"id":2}
2. 消息生产与消费
生产和消费均正常:
- 生产者命令:
kafka-json-schema-console-producer \ --bootstrap-server broker:9092 \ --property schema.registry.url=http://schema-registry:8081 \ --property value.schema.id=2 \ --property key.schema.id=1 \ --property key.separator='|' \ --property parse.key=true \ --topic locations {"profileId":"asdfghjkl"}|{"profileId":"asdfghjkl","latitude":90.000,"longitude":-180.000} {"profileId":"asdfghjkl"}|{"profileId":"asdfghjkl","latitude":90.000,"longitude":-179}
- 消费者命令:
kafka-json-schema-console-consumer \ --bootstrap-server broker:9092 \ --from-beginning \ --property schema.registry.url=http://localhost:8081 \ --property print.key=true \ --property key.separator='|' \ --topic locations {"profileId":"asdfghjkl"}|{"profileId":"asdfghjkl","latitude":90.000,"longitude":-180.000} {"profileId":"asdfghjkl"}|{"profileId":"asdfghjkl","latitude":90.000,"longitude":-179}
3. ksqlDB建表与查询错误
创建表:
ksql> CREATE TABLE loc WITH ( >KAFKA_TOPIC = 'locations', >KEY_FORMAT = 'JSON_SR', >KEY_SCHEMA_ID = 1, >VALUE_FORMAT = 'JSON_SR', >VALUE_SCHEMA_ID = 2 >);
表结构:
ksql> describe loc; Name : LOC Field | Type ------------------------------------------------------------- ROWKEY | STRUCT<profileId VARCHAR(STRING)> (primary key) profileId | VARCHAR(STRING) latitude | DOUBLE longitude | DOUBLE ------------------------------------------------------------- For runtime statistics and query details run: DESCRIBE <Stream,Table> EXTENDED;
执行查询时出现错误:
ksql> SELECT * FROM loc EMIT CHANGES; +-------------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------------------------+ |ROWKEY |profileId |latitude |longitude | +-------------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------------------------+ org.apache.kafka.streams.errors.StreamsException: Exception caught in process. taskId=0_0, processor=KSTREAM-SOURCE-0000000001, topic=locations, partition=0, offset=3, stacktrace=io.confluent.ksql.serde.KsqlSerializationException: Error serializing message to topic: _confluent-ksql-default_transient_transient_LOC_4102565114369480088_1698924612467-KsqlTopic-Reduce-changelog. Mismatching schema. Hint: You probably forgot to add VALUE_SCHEMA_ID when creating the source. at io.confluent.ksql.serde.connect.KsqlConnectSerializer.serialize(KsqlConnectSerializer.java:56) ... Caused by: org.apache.kafka.connect.errors.DataException: Mismatching schema. at io.confluent.connect.json.JsonSchemaData.fromConnectData(JsonSchemaData.java:524) ...
问题原因
- Schema的
$id冲突:手动为两个不同Schema指定了相同的$id值,而$id是Schema Registry用于唯一标识Schema的核心字段,手动设置会导致序列化/反序列化时的Schema匹配逻辑混乱。 - ksqlDB表定义缺少显式主键映射:创建表时未明确指定主键字段,仅依赖自动生成的
ROWKEY,导致ksqlDB在生成状态存储的changelog主题时,无法正确匹配序列化所需的Schema。
解决方法
1. 重新注册Schema(移除手动指定的$id)
修改Schema提交文件,删除手动设置的$id:
- 新的
schema-key-registry.json:
{ "schemaType":"JSON", "schema":"{\"$schema\":\"https://json-schema.org/draft/2020-12/schema\",\"title\":\"LocationKey\",\"type\":\"object\",\"properties\":{\"profileId\":{\"type\":\"string\"}}}" }
- 新的
schema-value-registry.json:
{ "schemaType":"JSON", "schema":"{\"$schema\":\"https://json-schema.org/draft/2020-12/schema\",\"title\":\"LocationValue\",\"type\":\"object\",\"properties\":{\"profileId\":{\"type\":\"string\",\"description\":\"The id of the location.\"},\"latitude\":{\"type\":\"number\",\"minimum\":-90,\"maximum\":90,\"description\":\"The location's latitude.\"},\"longitude\":{\"type\":\"number\",\"minimum\":-180,\"maximum\":180,\"description\":\"The location's longitude.\"}}}" }
重新注册(测试环境可先删除旧subject):
# 删除旧subject curl -X DELETE http://localhost:8081/subjects/locations-key curl -X DELETE http://localhost:8081/subjects/locations # 注册新Schema curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \ --data "$(cat schema-key-registry.json)" \ http://localhost:8081/subjects/locations-key/versions curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \ --data "$(cat schema-value-registry.json)" \ http://localhost:8081/subjects/locations/versions
2. 重新创建ksqlDB表(显式指定主键)
CREATE TABLE loc ( profileId VARCHAR PRIMARY KEY, latitude DOUBLE, longitude DOUBLE ) WITH ( KAFKA_TOPIC = 'locations', KEY_FORMAT = 'JSON_SR', VALUE_FORMAT = 'JSON_SR' );
若需指定Schema ID,替换为新注册返回的ID即可:
CREATE TABLE loc ( profileId VARCHAR PRIMARY KEY, latitude DOUBLE, longitude DOUBLE ) WITH ( KAFKA_TOPIC = 'locations', KEY_FORMAT = 'JSON_SR', KEY_SCHEMA_ID = <新的key schema ID>, VALUE_FORMAT = 'JSON_SR', VALUE_SCHEMA_ID = <新的value schema ID> );
3. 重新生产消息并查询
使用新的Schema ID重新生产消息,再执行查询:
SELECT * FROM loc EMIT CHANGES;
内容的提问来源于stack exchange,提问作者mrt181
相关产品推荐
相关产品推荐

