Skip to content

Flink SQL 行列转换:UNNEST 展开、LISTAGG 聚合与 JSON 数组处理

  1. 使用 UNNEST 将一行转换为多行
sql
-- 创建输入表
CREATE TABLE input_table (
  id INT,
  names ARRAY<STRING>
) WITH (
  ...
);

-- 使用 UNNEST 展开数组
SELECT id, name
FROM input_table CROSS JOIN UNNEST(names) AS T(name);

也可以用 Hive 的 explode 函数(需先加载 Hive 函数模块)

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

create table sourceTable (
    names ARRAY<VARCHAR>
) with (
    ...
);

create table sinkTable(
   name VARCHAR
) with (
    ...
);

insert into sinkTable select name from sourceTable, lateral table(explode(names)) as T(name);
  1. 使用 LISTAGG 将多行聚合成一行
sql
SELECT id, LISTAGG(value, ', ') AS aggregated_values
FROM your_table
GROUP BY id;
  1. 使用 collect_list 将多行聚合成一个 Array
sql
SELECT id, collect_list(value) AS aggregated_values
FROM your_table
GROUP BY id;
  1. 使用 collect_set 将多行聚合为一个 Array,并去重
sql
SELECT id, collect_set(value) AS aggregated_values
FROM your_table
GROUP BY id;
  1. JsonArray 字符串展开
sql
-- 用 UDF 把 JSON 数组字符串转成 ARRAY<STRING>
CREATE FUNCTION JsonArrayToList AS 'com.example.udf.JsonArrayToList';

CREATE TABLE input_table (
  id INT,
  names STRING
) WITH (
  'connector' = 'datagen'
);

SELECT id, name
FROM (
SELECT id, JsonArrayToList(names) as nameArray
FROM input_table) t1 CROSS JOIN UNNEST(nameArray) AS T(name);

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