如何让KSQLDB中Kafka流的键列仅显示JSON键值而非完整串?
解决KSQLDB中JSON键解析异常问题
问题根源
你的流定义将KEY_FORMAT设为JSON,但键字段id定义为VARCHAR类型。KSQL会把完整的JSON键对象(即{"id":"COLE888"})当作字符串直接赋值给id字段,而非自动解析JSON中的id属性值。
解决方案
方案一:调整流定义以匹配JSON键结构
修改流的键字段为STRUCT类型,对应JSON键的结构,后续查询时提取指定属性:
CREATE STREAM IF NOT EXISTS users_stream ( key STRUCT<id VARCHAR> KEY, employeeId BIGINT ) WITH ( KAFKA_TOPIC='users_topic', KEY_FORMAT='JSON', VALUE_FORMAT='JSON', PARTITIONS=2 );
查询时获取STRUCT中的id值:
SELECT key->id AS id, employeeId FROM users_stream;
方案二:改用字符串格式的键
如果不需要JSON格式的键,可将KEY_FORMAT改为KAFKA(对应字符串序列化),直接发送字符串形式的键:
- 删除原有流(若需重新定义):
DROP STREAM IF EXISTS users_stream;
- 重新创建流:
CREATE STREAM IF NOT EXISTS users_stream ( id VARCHAR KEY, employeeId BIGINT ) WITH ( KAFKA_TOPIC='users_topic', KEY_FORMAT='KAFKA', VALUE_FORMAT='JSON', PARTITIONS=2 );
- 修改生产者命令,使用字符串序列化器发送键:
docker-compose exec broker kafka-console-producer --broker-list localhost:29092 \ --property parse.key=true \ --property key.separator="&" \ --property key.serializer=org.apache.kafka.common.serialization.StringSerializer \ --property value.serializer=custom.class.serialization.JsonSerializer \ --topic users_topic
- 发送消息时使用纯字符串键:
COLE888&{"employeeId":"1470258"}
内容的提问来源于stack exchange,提问作者magda_lena
相关产品推荐
相关产品推荐

