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

ksqlDB无法反序列化JSON键:如何创建双主键源表?

问题描述

通过Kafka Connect的JDBC Source连接器将数据导入ksqlDB源表,已配置transform将c1、c2列转为Kafka主题的键,但创建以c1、c2为复合主键的源表时,ksqlDB无法反序列化键并报错,询问是否能实现这种双主键源表的创建。

连接器配置

CREATE SOURCE CONNECTOR IF NOT EXISTS loadshedding_model_source WITH (
  'connector.class'             = 'io.confluent.connect.jdbc.JdbcSourceConnector',
  'connection.url'              = 'jdbc:postgresql://cockroachdb:26257/db',
  'connection.user'             = 'user',
  'connection.password'         = 'password',
  'connection.options'          = '-c multiple_active_portals_enabled=true',
  'topic.prefix'                = 'jdbc_',
  'table.whitelist'             = 'whitelist_table',
  'mode'                        = 'bulk',
  'numeric.mapping'             = 'best_fit',
  "poll.interval.ms"            = 1800000,
  'transforms'                  = 'createKey',
  'transforms.createKey.type'   = 'org.apache.kafka.connect.transforms.ValueToKey',
  'transforms.createKey.fields' = 'c1,c2',
  'topic.creation.default.partitions' = 3,
  'topic.creation.default.replication.factor' = 3
);

预期主题消息键格式

{
  "c1": "c1v1",
  "c2": "v1"
}

创建源表语句

CREATE SOURCE TABLE whitelist_table (
  C1 VARCHAR PRIMARY KEY,
  C2 VARCHAR PRIMARY KEY,
  C3 VARCHAR,
  C4 TIMESTAMP,
  C5 INT
) WITH (
  KAFKA_TOPIC='jdbc_whitelist_table',
  KEY_FORMAT='JSON',
  VALUE_FORMAT='JSON'
)

错误信息

ERROR {"type":0,"deserializationError":{"target":"key","errorMessage":"Failed to deserialize key from topic: jdbc_model_load_shedding. Unrecognized token 'Struct': was expecting ( JSON String, Number, Array, Object or token 'null', 'true' or 'false')","recordB64":null,"cause":["Unrecognized token 'Struct': was expecting (JSON String, Number, Array, Object or token 'null', 'true' or 'false')"],"topic":"jdbc_whitelist_table"},"recordProcessingError":null,"productionError":null,"serializationError":null,"kafkaStreamsThreadError":null} (processing.CST_WHITELIST_TABLE_11.KsqlTopic.Source.deserializer:44)
解决方案

错误原因

ValueToKey转换生成的键是Kafka Connect Struct类型,而非JSON格式。你在ksqlDB中指定KEY_FORMAT='JSON',导致反序列化失败——实际键是Struct的字符串表示(比如Struct{c1=c1v1,c2=v1}),不是标准JSON结构。

修复步骤

  1. 修改连接器配置,添加Struct转JSON的transform
    在transforms中新增unwrapKey,将Struct类型的键序列化为标准JSON:

    CREATE SOURCE CONNECTOR IF NOT EXISTS loadshedding_model_source WITH (
      'connector.class'             = 'io.confluent.connect.jdbc.JdbcSourceConnector',
      'connection.url'              = 'jdbc:postgresql://cockroachdb:26257/db',
      'connection.user'             = 'user',
      'connection.password'         = 'password',
      'connection.options'          = '-c multiple_active_portals_enabled=true',
      'topic.prefix'                = 'jdbc_',
      'table.whitelist'             = 'whitelist_table',
      'mode'                        = 'bulk',
      'numeric.mapping'             = 'best_fit',
      "poll.interval.ms"            = 1800000,
      'transforms'                  = 'createKey,unwrapKey',  -- 新增unwrapKey转换
      'transforms.createKey.type'   = 'org.apache.kafka.connect.transforms.ValueToKey',
      'transforms.createKey.fields' = 'c1,c2',
      'transforms.unwrapKey.type'   = 'org.apache.kafka.connect.transforms.ExtractField$Key',
      'transforms.unwrapKey.field'  = '',  -- 空字段表示提取整个Struct为JSON
      'topic.creation.default.partitions' = 3,
      'topic.creation.default.replication.factor' = 3
    );
    
  2. 修正ksqlDB源表的主键定义
    ksqlDB中复合主键需用PRIMARY KEY (列1, 列2)的形式,不能给每个列单独标记PRIMARY KEY。修改后的创建语句:

    CREATE SOURCE TABLE whitelist_table (
      C1 VARCHAR,
      C2 VARCHAR,
      C3 VARCHAR,
      C4 TIMESTAMP,
      C5 INT,
      PRIMARY KEY (C1, C2)
    ) WITH (
      KAFKA_TOPIC='jdbc_whitelist_table',
      KEY_FORMAT='JSON',
      VALUE_FORMAT='JSON'
    );
    

验证

修改配置后重启连接器,确认Kafka主题的键为标准JSON格式,再执行创建表语句即可正常加载数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 13:35:31