如何使用Apache Flink SQL解码包含编码JSON的字符串
Flink 1.10 解析Filebeat嵌套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
相关产品推荐
相关产品推荐

