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

Apache Flink中为何不同Key会被分配至同一窗口?

问题:Flink窗口中出现不同Key事件的原因分析

输入事件代码

ctx.collect(new ClickEvent(2, 2, Instant.parse("2000-05-09T12:00:01.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(2, 2, Instant.parse("2000-05-09T12:00:02.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(2, 2, Instant.parse("2000-05-09T12:00:03.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(2, 2, Instant.parse("2000-05-09T12:00:06.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(2, 2, Instant.parse("2000-05-09T12:00:07.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(5, 5, Instant.parse("2000-05-09T12:00:07.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(5, 5, Instant.parse("2000-05-09T12:00:09.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(5, 5, Instant.parse("2000-05-09T12:00:10.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(2, 2, Instant.parse("2000-05-09T12:00:04.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(2, 2, Instant.parse("2000-05-09T12:00:05.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(5, 5, Instant.parse("2000-05-09T12:00:12.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(1, 1, Instant.parse("2000-05-09T12:00:13.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(1, 1, Instant.parse("2000-05-09T12:00:15.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(1, 1, Instant.parse("2000-05-09T12:00:16.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(5, 5, Instant.parse("2000-05-09T12:00:11.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(1, 1, Instant.parse("2000-05-09T12:00:16.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(1, 1, Instant.parse("2000-05-09T12:00:17.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(1, 1, Instant.parse("2000-05-09T12:00:15.00Z"), 1));
Thread.sleep(1500);
ctx.collect(new ClickEvent(1, 1, Instant.parse("2000-05-09T12:00:22.00Z"), 1));

Flink处理管道代码

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.getConfig().setAutoWatermarkInterval(1000);
SingleOutputStreamOperator<ClickEvent> output = env.addSource(source).assignTimestampsAndWatermarks
        (WatermarkStrategy.<ClickEvent>forBoundedOutOfOrderness(Duration.ofMillis(1000))
                .withTimestampAssigner((click, t) -> click.getEventTimeMillis())).keyBy((ClickEvent e) -> e.userId)
        .window(TumblingEventTimeWindows.of(Time.milliseconds(10000)))
        .process(new ProcessWindowFunction<ClickEvent, ClickEvent, Long, TimeWindow>() {
            @Override
            public void process(Long aLong, ProcessWindowFunction<ClickEvent, ClickEvent, Long, TimeWindow>.Context context, Iterable<ClickEvent> iterable, Collector<ClickEvent> collector) throws Exception {
                    System.out.println("[ Window Start: "+context.window().getStart()+" - Current Watermark: "+
                            context.currentWatermark() + " - Current Key: "+aLong + " - Window End: "+
                            context.window().getEnd() + " ]");
                    System.out.println("==== Window Elements ====");
                    for(ClickEvent clickEvent : iterable){
                        System.out.println(clickEvent.toString());
//                        collector.collect(clickEvent);
                    }
            }
        });
output.print();
return env.execute("Streaming Job with Global Watermark");

运行输出结果

[ Window Start: 957873600000 - Current Watermark: 957873610999 - Current Key: 2 - Window End: 957873610000 ]
==== Window Elements ====
ClickEvent{ clickId= 2, clickTimestamp= 2000-05-09T12:00:01Z, userId= 2, toSum=1}
ClickEvent{ clickId= 2, clickTimestamp= 2000-05-09T12:00:02Z, userId= 2, toSum=1}
ClickEvent{ clickId= 2, clickTimestamp= 2000-05-09T12:00:03Z, userId= 2, toSum=1}
ClickEvent{ clickId= 2, clickTimestamp= 2000-05-09T12:00:06Z, userId= 2, toSum=1}
ClickEvent{ clickId= 2, clickTimestamp= 2000-05-09T12:00:07Z, userId= 2, toSum=1}
ClickEvent{ clickId= 2, clickTimestamp= 2000-05-09T12:00:04Z, userId= 2, toSum=1}
ClickEvent{ clickId= 2, clickTimestamp= 2000-05-09T12:00:05Z, userId= 2, toSum=1}
[ Window Start: 957873600000 - Current Watermark: 957873610999 - Current Key: 5 - Window End: 957873610000 ]
==== Window Elements ====
ClickEvent{ clickId= 5, clickTimestamp= 2000-05-09T12:00:07Z, userId= 5, toSum=1}
ClickEvent{ clickId= 5, clickTimestamp= 2000-05-09T12:00:09Z, userId= 5, toSum=1}
[ Window Start: 957873610000 - Current Watermark: 9223372036854775807 - Current Key: 5 - Window End: 957873620000 ]
==== Window Elements ====
ClickEvent{ clickId= 5, clickTimestamp= 2000-05-09T12:00:10Z, userId= 5, toSum=1}
[ Window Start: 957873610000 - Current Watermark: 9223372036854775807 - Current Key: 1 - Window End: 957873620000 ]
==== Window Elements ====
ClickEvent{ clickId= 1, clickTimestamp= 2000-05-09T12:00:13Z, userId= 1, toSum=1}
ClickEvent{ clickId= 5, clickTimestamp= 2000-05-09T12:00:12Z, userId= 5, toSum=1}
ClickEvent{ clickId= 1, clickTimestamp= 2000-05-09T12:00:15Z, userId= 1, toSum=1}
ClickEvent{ clickId= 5, clickTimestamp= 2000-05-09T12:00:11Z, userId= 5, toSum=1}
ClickEvent{ clickId= 1, clickTimestamp= 2000-05-09T12:00:16Z, userId= 1, toSum=1}
ClickEvent{ clickId= 1, clickTimestamp= 2000-05-09T12:00:16Z, userId= 1, toSum=1}
ClickEvent{ clickId= 1, clickTimestamp= 2000-05-09T12:00:17Z, userId= 1, toSum=1}
ClickEvent{ clickId= 1, clickTimestamp= 2000-05-09T12:00:15Z, userId= 1, toSum=1}
[ Window Start: 957873620000 - Current Watermark: 9223372036854775807 - Current Key: 1 - Window End: 957873630000 ]
==== Window Elements ====
ClickEvent{ clickId= 1, clickTimestamp= 2000-05-09T12:00:22Z, userId= 1, toSum=1}

疑问

原本认为每个窗口应仅处理一个Key,但实际运行后却出现不同Key被分配至同一窗口的现象,想请教原因。


原因分析与解决

1. 同一时间窗口对应不同Key是正常现象

你看到的相同窗口起止时间但不同Key的日志是正常的:Flink中keyBy()之后的窗口是Keyed Window,每个Key会维护属于自己的独立窗口集合,所有Key的时间窗口划分是对齐的(比如10秒滚动窗口的时间区间是全局一致的),但每个Key的同时间区间窗口是完全独立的实例,彼此不共享数据。

比如日志中957873600000-957873610000窗口分别对应Key=2和Key=5,这是两个独立的窗口实例,各自处理对应Key的事件,属于正常逻辑。

2. 不同Key事件混入同一窗口实例的核心问题

而你遇到的Key=1的窗口里出现Key=5的事件,是因为Key类型不匹配导致的分组错误:

  • 你的keyBy((ClickEvent e) -> e.userId)返回的是Integer类型(假设userId是int字段),但ProcessWindowFunction的泛型Key声明为Long。
  • 这种类型不匹配会在运行时触发强制转换,破坏Flink的Key哈希分组逻辑,导致不同Key的事件被错误分配到同一个窗口实例中。

解决方法

统一Key的类型,保证keyBy返回值与ProcessWindowFunction的泛型Key一致:

  • 如果userId是int类型,修改ProcessWindowFunction的泛型为Integer:
.process(new ProcessWindowFunction<ClickEvent, ClickEvent, Integer, TimeWindow>() {
    @Override
    public void process(Integer userId, Context context, Iterable<ClickEvent> iterable, Collector<ClickEvent> collector) throws Exception {
        // 原有逻辑保持不变
    }
});
  • 或者将ClickEvent的userId字段改为Long类型,保持与泛型Key一致。

修改后,每个窗口实例只会处理对应Key的事件,不会出现跨Key的混流问题。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 09:05:33