如何使用Flink CEP检测同整数序列后接计数的模式?是否支持上下文传递?
使用Flink CEP检测连续相同整数序列+对应计数的模式
问题背景
需要从整数流中检测两类序列:
- 一段连续相同的整数序列,紧接着一个等于该序列长度的整数(即序列中整数的计数);
- 一段连续相同的整数序列,后续没有出现对应长度的整数(对应示例中的
List(1,1,1))。
示例输入:
val input = List(1,2,2,2,3,3,2,1,1,1,3,0,1)
期望匹配结果:
List(2,2,2,3) List(3,3,2) List(1,1,1)
希望通过Flink CEP实现,同时想了解CEP是否支持上下文传递,或能否通过修改事件传递额外数据。
实现方案
1. 预处理事件(推荐方式)
先对原始整数流做预处理,给每个事件添加当前连续相同值的长度字段,这样后续CEP模式可以直接基于这个字段匹配,逻辑更清晰。
定义包装类:
case class CountedEvent(value: Int, currentRunLength: Int)
预处理流计算连续长度:
val countedStream = inputStream .keyBy(_ => 1) // 全局单键,保证连续性判断覆盖全流 .process(new KeyedProcessFunction[Int, Int, CountedEvent] { private var lastVal: Int = _ private var runLen: Int = 0 override def processElement( value: Int, ctx: KeyedProcessFunction[Int, Int, CountedEvent]#Context, out: Collector[CountedEvent]): Unit = { if (value == lastVal) { runLen += 1 } else { runLen = 1 lastVal = value } out.collect(CountedEvent(value, runLen)) } })
2. 定义CEP模式
基于预处理后的CountedEvent流,定义两种模式:
模式1:匹配带计数的序列
匹配连续相同值序列,后续紧跟等于序列长度的整数:
import org.apache.flink.cep.scala.pattern.Pattern // 匹配连续相同值的起始和后续元素 val sameValueSeq = Pattern.begin[CountedEvent]("start") .where(_.currentRunLength == 1) // 标记新序列的起点 .followedByAny("continuation") .where(_.currentRunLength > 1) // 匹配序列的后续相同元素 .oneOrMore() .optional() // 允许序列只有1个元素 // 追加计数匹配条件 val withCountPattern = sameValueSeq .followedBy("count") .where((ctx, event) => { // 从上下文获取序列的最后一个元素,拿到连续长度 val lastSeqEvent = ctx.getEventsForPattern("continuation").lastOption .getOrElse(ctx.getEventsForPattern("start").head) event.value == lastSeqEvent.currentRunLength })
模式2:匹配无后续计数的序列
匹配连续相同值序列结束后,没有出现对应长度的整数的情况:
val withoutCountPattern = sameValueSeq .notFollowedBy("count") .where((ctx, event) => { val lastSeqEvent = ctx.getEventsForPattern("continuation").lastOption .getOrElse(ctx.getEventsForPattern("start").head) event.value == lastSeqEvent.currentRunLength })
3. 提取匹配结果
合并两种模式,提取原始整数序列:
import org.apache.flink.cep.scala.CEP import org.apache.flink.cep.scala.pattern.PatternStream val patternStream = CEP.pattern(countedStream, withCountPattern.union(withoutCountPattern)) val resultStream = patternStream.select { pattern => // 收集序列中的相同值事件 val seqEvents = pattern("start") ++ pattern.getOrElse("continuation", Iterable.empty) // 加入计数事件(如果存在) val allEvents = seqEvents ++ pattern.getOrElse("count", Iterable.empty) // 转换回原始整数列表 allEvents.map(_.value).toList }
关于上下文传递的说明
Flink CEP支持通过带上下文的条件函数传递匹配状态:
- 使用
where((ctx, event) => ...)的重载方法,ctx参数可以调用getEventsForPattern(patternName)获取之前模式分支匹配到的所有事件,以此实现上下文信息的传递。 - 除了通过预处理添加字段,也可以直接用
IterativeCondition在模式内部维护状态(比如记录当前序列的值和长度),但这种方式代码复杂度更高,预处理的方式更易维护。
内容的提问来源于stack exchange,提问作者Mohammad Esmaeilzadeh
相关产品推荐
相关产品推荐

