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';

更多推荐