Flink使用KeyedCoProcessFunction双流Join时如何延迟数ms输出未匹配数据到侧输出流
解决方案
核心实现思路是利用Flink KeyedCoProcessFunction 自带的定时器机制实现等待逻辑,全程无需引入窗口,对现有代码侵入性极低。
具体改造步骤:
- 新增两个托管状态:
- 待匹配点击记录存储:
ListState<Row>存储当前key下暂时未匹配到lookup数据的点击流记录 - 定时器标记状态:
ValueState<Long>存储当前key已注册的定时器触发时间,避免重复注册定时器
- 待匹配点击记录存储:
- 修改
processElement2逻辑:未匹配到lookup数据时,不直接发侧输出,而是将记录写入待匹配状态,若当前key无待触发的定时器,注册一个N毫秒后触发的定时器 - 修改
processElement1逻辑:lookup数据写入状态后,立刻检查是否有待匹配的点击记录,若有则直接完成join输出,清空待匹配状态并取消已注册的定时器 - 实现
onTimer方法:定时器触发时,重新检查所有待匹配记录,能匹配则join输出,不能匹配则发侧输出,最后清空对应状态
完整修改后代码
val joinStream = lookDataStream.keyBy(row -> row.<Long>getFieldAs("id")) .connect(clickDataStream.keyBy(row -> row.<Long>getFieldAs("lookupid"))) .process(new EnrichJoinFunction()); public static class EnrichJoinFunction extends KeyedCoProcessFunction<Long, Row, Row, Row> { final OutputTag<Row> outputTag = new OutputTag<Row>("side-output") {}; // 等待时长,可根据实际乱序情况调整,单位毫秒 private static final long WAIT_TIME_MS = 50; private MapState<Long, Row> lookupState = null; // 存储未匹配的点击流记录 private ListState<Row> pendingClickState = null; // 存储当前key注册的定时器时间,避免重复注册 private ValueState<Long> timerState = null; @Override public void open(Configuration parameters) throws Exception { val lookupStateDesc = new MapStateDescriptor<Long, Row>( "lookupState", TypeInformation.of(Long.class), TypeInformation.of(new TypeHint<Row>() {})); lookupStateDesc.enableTimeToLive(StateTtlConfig.newBuilder(Time.days(15)).build()); lookupState = getRuntimeContext().getMapState(lookupStateDesc); val pendingClickDesc = new ListStateDescriptor<>("pendingClickState", TypeInformation.of(new TypeHint<Row>() {})); // 待匹配状态设置较短的TTL,避免状态泄漏 pendingClickDesc.enableTimeToLive(StateTtlConfig.newBuilder(Time.minutes(1)).build()); pendingClickState = getRuntimeContext().getListState(pendingClickDesc); val timerStateDesc = new ValueStateDescriptor<>("timerState", Long.class); timerState = getRuntimeContext().getState(timerStateDesc); } @Override public void processElement1( Row lookupRow, KeyedCoProcessFunction<Long, Row, Row, Row>.Context ctx, Collector<Row> out) throws Exception { log.debug("Received Lookup Record" + RowUtils.printRow(lookupRow)); val id = lookupRow.<Long>getFieldAs("id"); if (!lookupState.contains(id)) { lookupState.put(id, lookupRow); } // lookup数据到达后,立刻处理待匹配的点击记录 Long registeredTimer = timerState.value(); if (registeredTimer != null) { // 取消已注册的定时器 ctx.timerService().deleteProcessingTimeTimer(registeredTimer); timerState.clear(); // 遍历所有待匹配的点击记录,直接join输出 Iterator<Row> pendingIterator = pendingClickState.get().iterator(); while (pendingIterator.hasNext()) { Row clickRow = pendingIterator.next(); val joinRow = join(clickRow, lookupRow); out.collect(joinRow); pendingIterator.remove(); } pendingClickState.clear(); } } @Override public void processElement2( Row clickRow, KeyedCoProcessFunction<Long, Row, Row, Row>.Context ctx, Collector<Row> out) throws Exception { log.debug("Received Click stream Record" + RowUtils.printRow(clickRow)); // 注意:你原代码中click流keyBy用的是lookupid字段,这里取id可能是笔误,建议核对 val id = clickRow.<Long>getFieldAs("lookupid"); if (lookupState.contains(id)) { // enrich join val joinRow = join(clickRow, lookupState.get(id)); out.collect(joinRow); } else { // 未匹配到则加入待匹配状态 pendingClickState.add(clickRow); // 没有注册过定时器则注册 if (timerState.value() == null) { long triggerTime = ctx.timerService().currentProcessingTime() + WAIT_TIME_MS; ctx.timerService().registerProcessingTimeTimer(triggerTime); timerState.update(triggerTime); } } } @Override public void onTimer(long timestamp, KeyedCoProcessFunction<Long, Row, Row, Row>.OnTimerContext ctx, Collector<Row> out) throws Exception { // 定时器触发,重新检查所有待匹配记录 Long id = ctx.getCurrentKey(); Iterator<Row> pendingIterator = pendingClickState.get().iterator(); while (pendingIterator.hasNext()) { Row clickRow = pendingIterator.next(); if (lookupState.contains(id)) { val joinRow = join(clickRow, lookupState.get(id)); out.collect(joinRow); } else { // 还是未匹配到,发送到侧输出 ctx.output(outputTag, clickRow); } pendingIterator.remove(); } // 清空状态 pendingClickState.clear(); timerState.clear(); } public Row join(Row clickRow, Row lookupRow) throws ParseException { Row joinedRow = new Row(RowKind.INSERT, 13); // row setter join ouput return joinedRow; } }
注意事项
- 等待时长
WAIT_TIME_MS可根据实际乱序情况调整,建议设为10~100毫秒,太长会增加状态压力和处理延迟 - 如果是事件时间场景,将上述处理时间定时器改为事件时间定时器即可,对应取
ctx.timestamp()加上等待时长注册 - 该方案同一key的多条待匹配记录会批量处理,性能开销极低,完全适配每秒1万条的点击流规模
内容的提问来源于stack exchange,提问作者gaurav miglani
相关产品推荐
相关产品推荐

