KSQLDB基于String类型值Topic创建表报错解决方案咨询
如何在KSQLDB中基于值为String类型的Kafka Topic创建表
场景说明
现有存储RDF数据的Kafka Topic,消息值为嵌入RDF内容的纯字符串格式。按照KSQLDB官方文档对匿名值场景的配置要求,需设置value_format='KAFKA'且WRAP_SINGLE_VALUE=false,最初编写的建表语句如下:
CREATE SOURCE TABLE source_table_proxy ( key VARCHAR PRIMARY KEY, value VARCHAR ) WITH ( KEY_FORMAT='KAFKA', VALUE_FORMAT='KAFKA', WRAP_SINGLE_VALUE=false, KAFKA_TOPIC = 'topic' );
所用Topic基础参数
- Key类型:STRING
- Value类型:STRING
- 分区数:12
- 副本数:1
报错信息
执行上述建表语句时抛出如下异常:
The 'KAFKA' format only supports a single field. Got: [`VALUE` STRING, `ROWPARTITION` INTEGER, `ROWOFFSET` BIGINT]
问题根因
KAFKA格式有严格约束:值侧仅支持存在1个字段。创建SOURCE源表时,KSQLDB会默认自动追加ROWPARTITION(存储消息所在分区号,整数类型)、ROWOFFSET(存储消息偏移量,长整数类型)两个系统伪列,和语句中显式定义的value字段合计3个字段,直接触发格式校验报错。
可行解决方法
方法1:关闭源表默认伪列自动生成特性
在执行建表语句前,先在会话级别设置配置项,关闭自动追加系统伪列的行为,原有建表语句即可正常运行:
-- 会话级关闭源表默认伪列生成 SET 'ksql.source.table.default.columns.enabled' = 'false'; -- 执行原建表语句 CREATE SOURCE TABLE source_table_proxy ( key VARCHAR PRIMARY KEY, value VARCHAR ) WITH ( KEY_FORMAT='KAFKA', VALUE_FORMAT='KAFKA', WRAP_SINGLE_VALUE=false, KAFKA_TOPIC = 'topic' );
该方案支持自定义值字段名称,建表完成后可直接通过value字段读取Topic中的原始RDF字符串内容。
方法2:使用默认值列,不显式声明自定义值字段
如果不修改会话配置,可以在建表时不手动定义值字段,KSQLDB会自动将KAFKA格式解析出的单值映射到系统默认列ROWVAL,建表语句如下:
CREATE SOURCE TABLE source_table_proxy ( key VARCHAR PRIMARY KEY ) WITH ( KEY_FORMAT='KAFKA', VALUE_FORMAT='KAFKA', WRAP_SINGLE_VALUE=false, KAFKA_TOPIC = 'topic' );
建表完成后,查询时引用ROWVAL即可获取原始字符串值,示例查询语句:
SELECT key, ROWVAL FROM source_table_proxy EMIT CHANGES;
注意:上述两种方案都不需要修改Topic本身的配置,也不需要额外做数据格式转换,可直接读取已存在的String类型Topic数据。
内容的提问来源于stack exchange,提问作者MaatDeamon
相关产品推荐
相关产品推荐

