You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用Flink SQL无需窗口/状态直接聚合DataStream<List<Row>>中的列表?

现有数据源及数据流定义如下:

// source
DataSource<List<Row>> source;

// listStream
SingleOutputStreamOperator<List<Row>> listRowStream = env
    .fromSource(source, WatermarkStrategy.noWatermarks(), "sourceName")
    .returns(Types.LIST(Types.ROW_NAMED(new String[]{"batch_id", "name", "value"}, Types.STRING, Types.STRING, Types.INT)));

数据流中的数据示例:

  • 数据示例1:
[
    Row["1", "a", 1],
    Row["1", "a", 2],
    Row["1", "b", 3]
]
  • 数据示例2:
[
    Row["2", "a", 4],
    Row["2", "b", 5],
    Row["2", "b", 6],
    Row["2", "c", 2]
]

之前尝试过两种方式,但都存在弊端:

  1. 使用WITH语句展开数组再通过窗口聚合,SQL大致为:
WITH exploded AS (SELECT UNNEST(list_col) AS row_col FROM source_table)
SELECT row_col.batch_id, row_col.name, SUM(row_col.value) 
FROM exploded 
GROUP BY row_col.batch_id, row_col.name, TUMBLE(...)
  1. 用FlatMap将数组转单条记录再窗口聚合,SQL大致为:
SELECT batch_id, name, SUM(value) 
FROM flatmapped_table 
GROUP BY batch_id, name, TUMBLE(...)

这两种方法的问题是窗口时间难设置:过大浪费资源,过小无法保证同一batch_id的数据落入同一窗口。

实际需求是:每个Kafka消息对应一个List,属于同一批次,希望无需窗口或Flink状态,直接对每个List独立做聚合计算。


解决方案

完全可以实现每个List独立聚合,不需要窗口或状态,核心思路是在单条消息的范围内完成聚合,避免跨消息的状态依赖。以下是两种可行方案:

先在DataStream阶段对每个List做本地聚合,再将聚合结果转为表供SQL使用:

// 对每个List独立聚合,得到每个batch_id+name的sum值
SingleOutputStreamOperator<Row> aggregatedStream = listRowStream
    .map(list -> {
        // 用HashMap做本地聚合
        Map<String, Integer> sumMap = new HashMap<>();
        for (Row row : list) {
            String batchId = row.getFieldAs("batch_id");
            String name = row.getFieldAs("name");
            int value = row.getFieldAs("value");
            String key = batchId + "_" + name;
            sumMap.put(key, sumMap.getOrDefault(key, 0) + value);
        }
        // 将聚合结果转为Row集合并展开
        return sumMap.entrySet().stream()
            .map(entry -> {
                String[] parts = entry.getKey().split("_");
                return Row.of(parts[0], parts[1], entry.getValue());
            })
            .collect(Collectors.toList());
    })
    .flatMap(Iterable::iterator)
    .returns(Types.ROW_NAMED(new String[]{"batch_id", "name", "sum_value"}, Types.STRING, Types.STRING, Types.INT));

// 转为Flink SQL表
tableEnv.createTemporaryView("aggregated_table", aggregatedStream);

// 直接用SQL查询,无需窗口
String sql = "SELECT batch_id, name, sum_value FROM aggregated_table";
Table resultTable = tableEnv.sqlQuery(sql);

通过UNNEST展开数组,结合消息唯一标识限定聚合范围:
首先将DataStream注册为表:

tableEnv.createTemporaryView("list_table", listRowStream, $("list_col"));

然后使用SQL完成单消息内的聚合:

SELECT 
    batch_id,
    name,
    SUM(value) AS sum_value
FROM (
    -- 展开每个消息的List,生成消息唯一标识限定聚合范围
    SELECT 
        UUID() AS msg_id,
        t.batch_id,
        t.name,
        t.value
    FROM list_table,
    UNNEST(list_col) AS t(batch_id, name, value)
)
GROUP BY msg_id, batch_id, name

msg_id是每条消息的唯一标识,GROUP BY msg_id确保聚合只在当前消息的范围内进行,不会跨消息,无需窗口或状态。

如果同一List内的batch_id完全一致(如示例中的情况),可简化SQL:

SELECT 
    batch_id,
    name,
    SUM(value) AS sum_value
FROM (
    SELECT 
        t.batch_id,
        t.name,
        t.value
    FROM list_table,
    UNNEST(list_col) AS t(batch_id, name, value)
)
GROUP BY batch_id, name

因为同一消息内的batch_id相同,GROUP BY batch_id, name自然只会在当前消息的范围内聚合。

关键说明

  • 两种方案均不依赖Flink状态,聚合计算完全在单条消息上下文内完成,避免了窗口带来的资源浪费和数据拆分问题。
  • 方案1适合Java代码优先的开发场景,方案2适合SQL优先的开发模式。

内容的提问来源于stack exchange,提问作者user29895709

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 03:27:28