如何使用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
相关产品推荐
相关产品推荐

