KSQL基于Kafka Topic建表后查询报key反序列化异常如何解决
问题根因
报错核心为Kafka消息Key反序列化类型不匹配:
你配置了
key_format='JSON',对应Topic的消息Key本身是JSON结构{"MESG_SEQ_NO":1015},但建表时直接将MESG_SEQ_NO声明为PRIMARY KEY,KSQL默认会尝试把整个Key的JSON对象转换为你指定的INT类型,类型不匹配导致反序列化失败,因此查询无返回结果。
解决方法
方案1:修改建表语句,明确主键来自Key的嵌套字段
KSQL中如果主键是JSON Key内的字段,需要明确标记字段来源,调整后的建表语句如下:
CREATE TABLE MSG_2639_SALES_TRNS_STORES_TABLE_S( MESG_SEQ_NO BIGINT PRIMARY KEY, MESSAGENO STRING, MESSAGECREATIONDATETIME TIMESTAMP, OPCO_GLN STRING, OPCO_COUNTRYCODE STRING, STORENO INT, SRVPNO INT, ENDDATETIMETRANSACTION TIMESTAMP, DATETIMESENTSTORE TIMESTAMP, MESG_ACTION STRING, MESG_IND_COMPLETED STRING, MESG_IND_SENTTOBROKER STRING ) WITH ( kafka_topic='ah_topics_2639_sales', key_format='JSON', value_format='JSON', -- 可选配置:需要查询历史数据时添加,从最早偏移量开始消费 'auto.offset.reset' = 'earliest' );
调整说明:
- 移除INT类型的长度参数,KSQL的数值类型不需要指定长度
- 将MESG_SEQ_NO类型调整为BIGINT,匹配报错要求的类型同时避免数值溢出
方案2:使用Value内的同名字段做主键
如果业务允许用Value中的MESG_SEQ_NO做主键,可以忽略Key结构,调整建表逻辑如下:
CREATE TABLE MSG_2639_SALES_TRNS_STORES_TABLE_S( MESG_SEQ_NO BIGINT PRIMARY KEY, MESSAGENO STRING, MESSAGECREATIONDATETIME TIMESTAMP, OPCO_GLN STRING, OPCO_COUNTRYCODE STRING, STORENO INT, SRVPNO INT, ENDDATETIMETRANSACTION TIMESTAMP, DATETIMESENTSTORE TIMESTAMP, MESG_ACTION STRING, MESG_IND_COMPLETED STRING, MESG_IND_SENTTOBROKER STRING ) WITH ( kafka_topic='ah_topics_2639_sales', value_format='JSON', -- 声明主键来自Value字段 PRIMARY_KEY_STRATEGY='VALUE', 'auto.offset.reset' = 'earliest' );
验证步骤
- 删除旧表:
DROP TABLE IF EXISTS MSG_2639_SALES_TRNS_STORES_TABLE_S; - 执行调整后的建表语句
- 重新执行查询:
select * from MSG_2639_SALES_TRNS_STORES_TABLE_S emit changes;即可正常返回数据
内容的提问来源于stack exchange,提问作者sridhar reddy
相关产品推荐
相关产品推荐

