如何在Flink DDL中解析Kafka主题内的两种不同JSON根消息?
解析同一Kafka主题中两种根结构JSON消息的Flink DDL方案
问题核心
同一Kafka主题下存在两种根结构不同的JSON消息({'Message1': {'b': 'c'}}和{'Message2': {'e': 'f'}}),直接用固定结构的JSON格式解析会失败,需要通过灵活方式兼容两种结构。
方案一:读取原始字符串后用JSON函数解析
这种方式先将Kafka消息的value以原始字符串形式读取,再通过Flink的JSON函数区分并解析两种结构,兼容性最强。
1. 修改Kafka源表DDL
将message字段改为STRING类型,并用raw格式读取原始JSON字符串,避免因结构不匹配导致解析报错:
CREATE TABLE audienceInput ( `messageKey` VARBINARY, `message` STRING, `topic` VARCHAR ) WITH ( 'connector' = 'kafka', 'topic' = 'mytopic', 'properties.bootstrap.servers' = '****:9092', 'scan.startup.mode' = 'earliest-offset', 'value.format' = 'raw' );
2. 解析两种消息结构
通过JSON_VALUE函数判断根字段是否存在,分别提取对应内容:
-- 创建视图统一解析结果 CREATE VIEW parsedMessages AS SELECT messageKey, topic, -- 提取具体业务字段 CASE WHEN JSON_VALUE(message, '$.Message1') IS NOT NULL THEN JSON_VALUE(message, '$.Message1.b') WHEN JSON_VALUE(message, '$.Message2') IS NOT NULL THEN JSON_VALUE(message, '$.Message2.e') ELSE NULL END AS content, -- 标记消息类型 CASE WHEN JSON_VALUE(message, '$.Message1') IS NOT NULL THEN 'Message1' WHEN JSON_VALUE(message, '$.Message2') IS NOT NULL THEN 'Message2' ELSE 'Unknown' END AS messageType FROM audienceInput;
如果需要完整解析嵌套结构,也可以用JSON_OBJECT构造结构化数据:
SELECT messageKey, topic, CASE WHEN JSON_VALUE(message, '$.Message1') IS NOT NULL THEN JSON_OBJECT('b' VALUE JSON_VALUE(message, '$.Message1.b')) WHEN JSON_VALUE(message, '$.Message2') IS NOT NULL THEN JSON_OBJECT('e' VALUE JSON_VALUE(message, '$.Message2.e')) END AS parsedContent FROM audienceInput;
方案二:用嵌套MAP类型兼容两种结构
利用Flink的MAP类型接收任意JSON结构,配合json.fail-on-missing-field参数避免解析报错。
1. 修改Kafka源表DDL
将message定义为嵌套MAP类型,允许缺失字段:
CREATE TABLE audienceInput ( `messageKey` VARBINARY, `message` MAP<STRING, MAP<STRING, STRING>>, `topic` VARCHAR ) WITH ( 'connector' = 'kafka', 'topic' = 'mytopic', 'properties.bootstrap.servers' = '****:9092', 'scan.startup.mode' = 'earliest-offset', 'value.format' = 'json', 'json.fail-on-missing-field' = 'false' );
2. 查询解析结果
通过MAP的key判断消息类型,提取对应字段:
SELECT messageKey, topic, CASE WHEN message['Message1'] IS NOT NULL THEN message['Message1']['b'] WHEN message['Message2'] IS NOT NULL THEN message['Message2']['e'] ELSE NULL END AS content FROM audienceInput;
内容的提问来源于stack exchange,提问作者elademir
相关产品推荐
相关产品推荐

