我们如何使用 Flink SQL API 定义嵌套的 json 属性(包括数组)?

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)