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

如何使用Flink CEP检测同整数序列后接计数的模式?是否支持上下文传递?

问题背景

需要从整数流中检测两类序列:

  1. 一段连续相同的整数序列,紧接着一个等于该序列长度的整数(即序列中整数的计数);
  2. 一段连续相同的整数序列,后续没有出现对应长度的整数(对应示例中的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 15:20:51