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
相关产品推荐
相关产品推荐

