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

使用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)
        ...

问题原因

  1. Schema的$id冲突:手动为两个不同Schema指定了相同的$id值,而$id是Schema Registry用于唯一标识Schema的核心字段,手动设置会导致序列化/反序列化时的Schema匹配逻辑混乱。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 08:44:49