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

如何使用Apache Flink SQL解码包含编码JSON的字符串

不需要自定义UDF,Flink 1.10已经内置了JSON解析相关功能,可以直接用两种方式实现需求:

方案1:用内置JSON函数动态解析

先在Kafka表定义中将外层Filebeat的JSON结构解析出来,message字段先定义为STRING类型,后续查询时调用内置函数提取内部属性:

-- 第一步:创建Kafka源表,先解析最外层Filebeat JSON结构
CREATE TABLE kafka_filebeat_log (
  `@timestamp` STRING,
  `message` STRING -- 存储Filebeat封装的原始日志JSON字符串
) WITH (
  'connector' = 'kafka',
  'topic' = '你的Kafka Topic名称',
  'properties.bootstrap.servers' = '你的Kafka集群地址',
  'format' = 'json',
  'scan.startup.mode' = 'earliest-offset'
);

查询时调用JSON_VALUE提取普通字段、JSON_QUERY提取复杂对象/数组即可:

SELECT
  `@timestamp` AS collect_time,
  JSON_VALUE(`message`, '$.someProperty') AS some_property,
  JSON_QUERY(`message`, '$.complexObj') AS complex_object
FROM kafka_filebeat_log;

方案2:预定义ROW结构自动解析

如果message内的JSON结构固定,可以直接将message字段定义为和内部JSON匹配的ROW类型,Flink会自动完成解析,性能比动态调用函数更好:

CREATE TABLE kafka_filebeat_log (
  `@timestamp` STRING,
  `message` ROW<
    someProperty STRING,
    otherNum INT,
    nestedObj ROW<subField STRING>
  > -- 和原始日志JSON结构完全匹配的ROW定义
) WITH (
  'connector' = 'kafka',
  'topic' = '你的Kafka Topic名称',
  'properties.bootstrap.servers' = '你的Kafka集群地址',
  'format' = 'json',
  'scan.startup.mode' = 'earliest-offset'
);

该方式可以直接访问嵌套属性:

SELECT message.someProperty, message.nestedObj.subField FROM kafka_filebeat_log;

注意事项

  • 只有当message内的JSON结构固定时推荐用方案2,结构不固定的场景优先用方案1的JSON函数动态取值。
  • 如果有非常规JSON解析需求,内置函数无法满足时再考虑自定义UDF实现。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 01:36:00