【发布时间】:2020-10-21 02:40:08
【问题描述】:
我正在使用 Flink CEP 来检测来自 Kafka 的事件的模式。为简单起见,事件只有一种类型。我正在尝试检测连续事件流中字段值的变化。代码如下所示
val streamEnv = StreamExecutionEnvironment.getExecutionEnvironment
streamEnv.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)
streamEnv.addSource(new FlinkKafkaConsumer[..]())
.filter(...)
.map(...)
.assignTimestampsAndWatermarks(
WatermarkStrategy.forMonotonousTimestamps[Event]().withTimestampAssigner(..)
)
.keyBy(...)(TypeInformation.of(classOf[...]))
val pattern: Pattern[Event, _] =
Pattern.begin[Event]("start", AfterMatchSkipStrategy.skipPastLastEvent()).times(1)
.next("middle")
.oneOrMore()
.optional()
.where(new IterativeCondition[Event] {
override def filter(event: Event, ctx:...): Boolean = {
val startTrafficEvent = ctx.getEventsForPattern("start").iterator().next()
startTrafficEvent.getFieldValue().equals(event.getFieldValue())
}
})
.next("end").times(1)
.where(new IterativeCondition[Event] {
override def filter(event: Event, ctx:...): Boolean = {
val startTrafficEvent = ctx.getEventsForPattern("start").iterator().next()
!startTrafficEvent.getFieldValue().equals(event.getFieldValue())
}
})
.within(Time.seconds(30))
Kafka 主题有 104 个分区,事件均匀分布在各个分区上。当我提交作业时,parallelism 设置为 104。
在 Web UI 中,有 2 个任务:第一个是 Source->filter->map->timestamp/watermark;第二个是CepOperator->sink。每个任务有 104 个并行度。
子任务的工作量不均衡,应该来自keyBy。子任务之间的水印不一样,但是开始卡在一个数值上,很长一段时间没有变化。从日志中,我可以看到 CEP 不断评估事件,并将匹配的结果推送到下游接收器。
事件速率为 10k/s,第一个任务的背压保持high,第二个任务保持ok。
请帮助解释 CEP 中发生了什么以及如何解决该问题
谢谢
【问题讨论】:
标签: apache-flink flink-streaming flink-cep flink-sql