mri*_*cat 5 apache-flink flink-streaming flink-sql
我们在使用 Flink SQL 时遇到以下问题:我们已经配置了 Kafka Twitter 连接器以将推文添加到 Kafka,并且我们希望使用 Flink SQL 从表中的 Kafka 读取推文。
我们如何使用 Flink SQL API 定义嵌套的 json 属性(包括数组)?
我们尝试了以下方法,但不起作用(返回的值为空):
CREATE TABLE kafka_tweets(
payload ROW(`HashtagEntities` ARRAY[VARCHAR])
) WITH (
'connector' = 'kafka',
'topic' = 'twitter_status',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'scan.startup.mode' = 'earliest-offset',
'format' = 'json'
)
Run Code Online (Sandbox Code Playgroud)
在 Twitter 响应中,HashtagEntities 是一个对象数组。
小智 0
CREATE TABLE `table` (
`userid` BIGINT,
`json_data` VARCHAR(2147483647),
`request_id` AS JSON_VALUE(`json_data`, '$.request_id'),
`items` ARRAY<ROW<`itemid` BIGINT, `shopid` BIGINT>>,
`event_time` AS `TO_TIMESTAMP`(`FROM_UNIXTIME`(`timestamp`, 'yyyy-MM-dd HH:mm:ss')),
`version` AS `TO_TIMESTAMP`(`FROM_UNIXTIME`(`timestamp`, 'yyyy-MM-dd')),
WATERMARK FOR `event_time` AS `event_time` - INTERVAL '1' MINUTE
)
Run Code Online (Sandbox Code Playgroud)
| 归档时间: |
|
| 查看次数: |
895 次 |
| 最近记录: |