主题
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类型匹配表:
| Thrift | FlinkSQL |
|---|---|
| struct | Row |
| list | Array |
| time | Bigint |
| set | Array |
| enum | String |
嵌套格式使用
嵌套字段的引用
以上文的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。