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

Flink SQL 行列转换:UNNEST 展开、LISTAGG 聚合与 JSON 数组处理 「我的数据空间」实时计算实践笔记 · Flink SQL 系列Flink SQL 行列转换使用 UNNEST 将一行转换为多行-- 创建输入表CREATETABLEinput_table(idINT,names ARRAYSTRING)WITH(...);-- 使用 UNNEST 展开数组SELECTid,nameFROMinput_tableCROSSJOINUNNEST(names)AST(name);也可以用 Hive 的 explode 函数需先加载 Hive 函数模块-- 加载 Hive 函数模块后才能用 explode / collect_list 等 Hive 内置函数LOADMODULE hiveWITH(hive-version3.1.3);createtablesourceTable(names ARRAYVARCHAR)with(...);createtablesinkTable(nameVARCHAR)with(...);insertintosinkTableselectnamefromsourceTable,lateraltable(explode(names))asT(name);使用 LISTAGG 将多行聚合成一行SELECTid,LISTAGG(value,, )ASaggregated_valuesFROMyour_tableGROUPBYid;使用 collect_list 将多行聚合成一个 ArraySELECTid,collect_list(value)ASaggregated_valuesFROMyour_tableGROUPBYid;使用 collect_set 将多行聚合为一个 Array并去重SELECTid,collect_set(value)ASaggregated_valuesFROMyour_tableGROUPBYid;JsonArray 字符串展开-- 用 UDF 把 JSON 数组字符串转成 ARRAYSTRINGCREATEFUNCTIONJsonArrayToListAScom.example.udf.JsonArrayToList;CREATETABLEinput_table(idINT,names STRING)WITH(connectordatagen);SELECTid,nameFROM(SELECTid,JsonArrayToList(names)asnameArrayFROMinput_table)t1CROSSJOINUNNEST(nameArray)AST(name);本文收录于「我的数据空间」技术库——一套可私有化部署的数据平台(数据集成 / 实时计算 / 数据湖 / 湖仓查询 / 智能问数)。产品介绍见我的数据空间官网,支持私有化部署与 OEM 合作。