Flink训练练习CEP示例疑似Bug:LongRides逻辑与目标不符
关于Flink CEP识别未完成行程的疑问解答
你完全没错!这段代码里的completedRides模式确实是在匹配120分钟内完成的行程——先捕获一个标记为START的行程事件,然后在接下来的2小时内捕获对应的END事件。这看起来和练习目标「输出前2小时内未匹配到END事件的行程START事件」正好相反,但其实关键在于如何利用Flink CEP的超时处理机制来获取我们真正需要的结果。
Flink CEP中,当你给模式添加.within(Time.minutes(120))限制后,所有在120分钟窗口内没有完成后续匹配(也就是没找到对应END事件)的START事件,会被标记为超时模式实例。我们的目标就是从这些超时实例里提取出最初的START事件。
下面是实现练习目标的完整代码示例,基于你给出的模式进行扩展:
import org.apache.flink.cep.PatternSelectFunction; import org.apache.flink.cep.PatternTimeoutFunction; import org.apache.flink.cep.pattern.Pattern; import org.apache.flink.cep.pattern.conditions.SimpleCondition; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator; import org.apache.flink.streaming.api.windowing.time.Time; import org.apache.flink.util.OutputTag; // 定义匹配已完成行程的模式 Pattern<TaxiRide, TaxiRide> completedRides = Pattern.<TaxiRide>begin("start") .where(new SimpleCondition<TaxiRide>() { @Override public boolean filter(TaxiRide ride) throws Exception { return ride.isStart; } }) .next("end") .where(new SimpleCondition<TaxiRide>() { @Override public boolean filter(TaxiRide ride) throws Exception { return !ride.isStart; } }) .within(Time.minutes(120)); PatternStream<TaxiRide> patternStream = CEP.pattern(keyedRides, completedRides); // 创建侧输出标签,用于收集超时的未完成行程 OutputTag<TaxiRide> longRidesTag = new OutputTag<TaxiRide>("long-rides"){}; // 同时处理匹配成功的行程和超时的未完成行程 SingleOutputStreamOperator<TaxiRide> processedStream = patternStream.select( longRidesTag, // 处理超时实例:提取对应的START事件 (PatternTimeoutFunction<TaxiRide, TaxiRide>) (timeoutPattern, timeoutTimestamp) -> { return timeoutPattern.get("start").get(0); }, // 处理匹配成功的已完成行程:这里我们不需要,返回null即可 (PatternSelectFunction<TaxiRide, TaxiRide>) matchedPattern -> null ); // 从侧输出流中获取2小时内未完成的行程START事件 DataStream<TaxiRide> longRides = processedStream.getSideOutput(longRidesTag);
关键解释:
- 我们使用
select()方法的重载版本,它允许同时定义超时处理器和匹配成功处理器; - 超时处理器会接收所有在120分钟内没找到对应END事件的START事件,我们只需要从模式映射中取出
start对应的事件即可; - 匹配成功的已完成行程不是我们的目标,所以这里返回null忽略掉,最终我们需要的结果都在侧输出流
longRides里。
简单来说,你看到的初始代码只是完成了模式定义的第一步,练习的核心是通过处理这些超时的模式实例,反向得到未完成的行程数据。你的观察非常敏锐,一下子就抓住了模式本身的逻辑~
内容的提问来源于stack exchange,提问作者YuFeng Shen
相关产品推荐
相关产品推荐

