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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 20:09:03