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

如何使用ksqldb从JSON字符串payload中动态提取字段并序列化输出

实现方案

步骤1:创建原始CDC数据流

首先创建ksqldb流映射MongoDB Source Connector写入的原始Kafka Topic,外层数据本身是标准JSON格式,只需提取payload字符串字段即可:

CREATE STREAM mongo_cdc_raw (
  payload STRING
) WITH (
  KAFKA_TOPIC='替换为你的MongoDB源Topic名称',
  VALUE_FORMAT='JSON',
  PARTITIONS=替换为对应Topic的分区数
);

步骤2:解析嵌套JSON并生成目标数据

使用ksqldb内置的JSON_PARSE函数将payload字符串转为结构化JSON对象,再通过->操作符提取嵌套字段,直接生成符合要求的输出流:

CREATE STREAM parsed_mongo_comments WITH (
  KAFKA_TOPIC='parsed_comments_topic', -- 自定义输出的Kafka Topic名
  VALUE_FORMAT='JSON' -- 输出为标准JSON格式
) AS
SELECT
  JSON_PARSE(payload)->'fullDocument'->_id->'$oid' AS id,
  JSON_PARSE(payload)->'fullDocument'->'name' AS name,
  JSON_PARSE(payload)->'fullDocument'->'email' AS email,
  JSON_PARSE(payload)->'fullDocument'->'movie_id'->'$oid' AS movie_id,
  CAST(JSON_PARSE(payload)->'fullDocument'->'date'->'$date' AS BIGINT) AS date
FROM mongo_cdc_raw
EMIT CHANGES;

补充说明

  • 若需要输出fullDocument下的全部字段无需逐个提取,可直接返回JSON_PARSE(payload)->'fullDocument'作为值,ksqldb会自动序列化为标准JSON结构
  • 若payload中fullDocument的字段不固定需要完全动态序列化,可直接使用TO_JSON_STRING(JSON_PARSE(payload)->'fullDocument')输出整段JSON,无需逐个定义字段,可适配字段动态变化的场景
  • 若需要做类型校验,可对提取的字段增加CAST转换为指定类型,避免后续计算出现类型异常
  • 该方案兼容ksqldb 0.15及以上所有版本

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 03:06:03