Skip to content

FlinkSQL 嵌套类型处理:ROW / MAP / ARRAY 的声明、引用与构造

本文主要介绍在使用FlinkSQL时, 对于一些嵌套格式(MAP, ARRAY, STRUCT) 的一些处理方式

嵌套格式的声明

以一个用户行为结构体为例: 里面包含了list,嵌套结构等等, 这里不做详细展开:

c
struct ExtraAttribute {
        1:optional string originDeviceIds; // 原始设备号列表
    2:optional string appId;        // 应用 ID
    3:optional string clientIp;        // 客户端IP
    4:optional string msgId;        // 消息 ID
    5:optional string title;        // 消息标题
    6:optional string uuid;        // UUID
    .....
}
struct UserBehavior {
     1:required string id;// ID
     2:required BehaviorIdType id_type;// ID类型. 参考枚举 BehaviorIdType
     3:required DataSource data_source;// 数据源. 参考枚举 DataSource
     4:required ActionType action; // 用户行为. 参考枚举 ActionType
     5:optional string content;        // 行为内容.一般为文本,App为app_name,push为message等
     6:optional list<string> tags;        // 标签,一般是对文本,app做的标签抽取.
     7:optional ExtraAttribute ext;        // 非共用属性字段.每个数据源甚至行为可能有不同的外部属性
     8:required i64 time;        // 行为发生时间, 13位时间戳
}

对应的FlinkSQL的DDL声明如下所示: ExtraAttribute是嵌套的struct, 对应FlinkSQL的类型为Row<xxx xxx, xxx xxx>; id_type BehaviorIdType为Enum类型,对应FlinkSQL中的类型为String

sql
CREATE TABLE HDFSource (
        ext ROW<`originDeviceIds` STRING, `appId` STRING, `clientIp` STRING, `msgId` STRING, `title` STRING,  ....>,
        `action` STRING,
        id_type STRING,
        id STRING,
        `time` BIGINT,
        data_source STRING,
        content STRING,
        tags ARRAY<STRING>
) with (
        'connector.path'='hdfs://nameservice1/user/hive/warehouse/demo/data',
        'connector.type'='hdfs',
        'format.type'='parquet',
        'format.derive-schema'='true'
);

对于比较复杂的结构,建议在Kafka-platform绑定schema之后,由平台自动生成DDL.

Thrift -> FlinkSQL类型匹配表:

ThriftFlinkSQL
structRow
listArray
timeBigint
setArray
enumString

嵌套格式使用

嵌套字段的引用

以上文的ext的嵌套格式为例: 使用字段名.嵌套字段名即可引用

sql
select ext.originDeviceIds as deviceIds, ext.appId as appId from KafkaSource;

数组, map的取值

数组使用下标取值,map使用key来取值.

sql
create table KafkaSource (
  testMap MAP<STRING, STRING>, 
  testArray ARRAY<STRING>
) with (
  'connector.type' = 'xxx';
);

select testMap['name'] as name, testMap['age'] as age from KafkaSource;
select testArray[1] as name, testArray[2] as age from KafkaSource;

嵌套字段的构造

Flink提供了ROW, MAP, ARRAY 三个内置函数分别用于构造Row, MAP, ARRAY.

sql
create table  KafkaSink (
  testRow ROW<name STRING, age STRING>,
  testMap Map<STRING, STRING>,
  testArray ARRAY<STRING>
);

create table KafkaSource(
  name STRING,
  age STRING
);

insert into KafkaSink 
select 
ROW(name, age) as testRow,
MAP['name', name, 'age', age] as testMap,
ARRAY[name, age] as testArray
from KafkaSource;

数组的展开和构造

在数据的转换过程中,可能会存在将一行数组展开为多行, 也会存在将单行字段聚合组成一个数组或者list.

数组的展开, 一行展开为多行

借助 Hive 内置函数 explode 展开。

sql
-- 加载 Hive 函数模块后才能用 explode / collect_list 等 Hive 内置函数
LOAD MODULE hive WITH ('hive-version' = '3.1.3');

create table KafkaSource (
  id VARCHAR,
  bookList ARRAY<VARCHAR>
) with (
   ...
);

create table KafkaSink (
  id VARCHAR,
  book VARCHAR
) with (
   ...
);

insert into KafkaSink 
select id, book
from KafkaSource, lateral table(explode(bookList)) as T(book);

数组的构造, 多行聚合为一行

借助 Hive 内置函数 collect_list 聚合

sql
-- 加载 Hive 函数模块后才能用 explode / collect_list 等 Hive 内置函数
LOAD MODULE hive WITH ('hive-version' = '3.1.3');

create table KafkaSource (
  id VARCHAR,
  book VARCHAR
) with (
   ...
);

create table KafkaSink (
  id VARCHAR,
  bookList ARRAY<VARCHAR>
) with (
   ...
);

insert into KafkaSink
select id, collect_list(book)
from KafkaSource
group by id;

本文是「我的数据空间」实时计算实践笔记的一篇。平台把这类链路做成了托管作业(Flink SQL 流作业、实时集成、Iceberg 表管理与维护),支持私有化部署与 OEM 合作 —— 见核心能力总览技术博客,交流合作 QQ:1559851993。