Flink中如何使用multiset?怎样获取multiset中的元素?
Flink SQL multiset类型元素获取方案
Flink SQL原生对multiset类型的内置函数支持较少,你可以通过以下两种通用方式处理:
方案1:转换为ARRAY类型后操作
将multiset强制转换为同元素类型的ARRAY,即可直接使用所有数组相关函数,包括你提到的CARDINALITY取元素个数、下标取值等操作:
- 示例代码:
-- 假设collect聚合的是INT类型的user_id字段,分组字段为group_id SELECT group_id, CAST(collect(user_id) AS ARRAY<INT>) AS user_id_arr, CARDINALITY(CAST(collect(user_id) AS ARRAY<INT>)) AS user_count, CAST(collect(user_id) AS ARRAY<INT>)[1] AS first_user_id FROM source_table GROUP BY group_id
- 兼容适配:如果使用的是Flink 1.13及更早版本,直接CAST可能报错,可以先转JSON再解析为数组:
SELECT JSON_VALUE(CAST(collect(user_id) AS STRING), '$' RETURNING ARRAY<INT>) AS user_id_arr FROM source_table GROUP BY group_id
方案2:直接用UNNEST拆分为单行
如果需要把multiset中的每个元素展开为独立的行进行后续处理,可以直接用UNNEST函数,不需要提前做类型转换:
- 示例代码:
SELECT t.group_id, user_id FROM ( SELECT group_id, collect(user_id) AS user_ids FROM source_table GROUP BY group_id ) t, LATERAL TABLE(UNNEST(user_ids)) AS u(user_id)
小提示:如果是Flink 1.17及以上版本,官方已经新增了部分multiset专属操作函数,比如
MULTISET_CARDINALITY、MULTISET_CONTAINS等,可以直接调用无需转换。
内容的提问来源于stack exchange,提问作者volity
相关产品推荐
相关产品推荐

