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

如何让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(对应字符串序列化),直接发送字符串形式的键:

  1. 删除原有流(若需重新定义):
DROP STREAM IF EXISTS users_stream;
  1. 重新创建流:
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
);
  1. 修改生产者命令,使用字符串序列化器发送键:
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
  1. 发送消息时使用纯字符串键:
COLE888&{"employeeId":"1470258"}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 12:05:33