Flink解析Kafka中的复杂json
·
json格式如下:
{
"data": {
"Type": "xxx",
"timestamp": 1761013490373,
"logid": "xxx",
"info": {
"eventId": "xxx"
},
"baseparams": {
"osVersion": "0"
}
}
}
Flink代码:
-- 创建Kafka源表
CREATE TABLE kafka_source (
`data` ROW<
`Type` STRING,
`timestamp` BIGINT,
`logid` STRING,
`info` ROW<
`eventId` STRING
>
>,
`baseparams` ROW<
`osVersion` STRING
>
>
) WITH (
'connector' = 'kafka',
'topic' = 'sg',
'properties.bootstrap.servers' = 'example1:9092,example2:9092,example3:9092', -- 替换为实际Kafka地址
'properties.group.id' = 'App', -- 消费者组名
'format' = 'json',
'scan.startup.mode' = 'latest-offset',
'json.ignore-parse-errors' = 'true'
);
-- 创建StarRocks目标表
CREATE TABLE starrocks_sink (
`event_id` STRING,
`time_stamp` TIMESTAMP(3),
`os_version` STRING
) WITH (
'connector' = 'starrocks',
-- 支持传入多个地址,使用英文逗号 (,) 分隔
'jdbc-url' = 'jdbc:mysql:loadbalance://192.168.0.1:9030,192.168.0.2:9030,192.168.0.3:9030',
-- 支持传入多个地址,使用英文分号 (;) 分隔
'load-url' = '192.168.0.1:8030;192.168.0.1:8030;192.168.0.3:8030',
'database-name' = 'test',
'table-name' = 'demo_sr',
'username' = 'test',
'password' = '123456'
);
-- 插入数据到StarRocks
INSERT INTO starrocks_sink
SELECT
data.info.eventId AS event_id,
TO_TIMESTAMP(FROM_UNIXTIME(data.`timestamp` / 1000 )) AS time_stamp,
data.baseparams.osVersion AS os_version
FROM kafka_source
-- 可根据需要添加过滤条件
WHERE data.info.eventId = 'others';
更多推荐
所有评论(0)