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

如何在Flink DDL中解析Kafka主题内的两种不同JSON根消息?

问题核心

同一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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 22:15:41