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结构。
修复步骤
修改连接器配置,添加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 );修正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
相关产品推荐
相关产品推荐

