【发布时间】:2020-03-06 01:57:45
【问题描述】:
在过去几天研究 Flink CEP 库时,我的印象是它并没有在 Flink 的标准功能中添加任何新的基础功能。看起来 Flink CEP 的唯一目的是让事件处理更容易,具有清晰的语义和直观的代码结构。例如,Flink CEP 仅呈现 5 semantics 的事件匹配跳过。虽然这些语义对于大范围的情况可能已经足够了,但它可能无法解决具体的问题,这让我们回到了普通的 Flink。
测试用例是以下模式:
Emmit a alert(represented by 'a') for each non-overlapping pair of numbers in a stream
由模式表示:
Pattern.begin[EventType]("pair",skipStrategy).where(new AlwaysTrueFunction()).times(2)
因此,对于像(在流中从左到右输入的数字)1 1 1 1 1 这样的输入,预期的输出将是 a a,但 5 种匹配跳过策略都不会给出正确的结果:
No-skip: a a a a
Skip-to-next: a a a a
Skip-past-last-event: a a a a
Skip-to-first[1]: a a a a
Skip-to-last[1]: a a a a
尽管这些策略无法生成所需的模式,但可以使用RichFunction 和ValueState 计数器轻松确定何时应发出新警报,将输入流转换为事件流。
因此,我希望对这些问题有所了解:
如果 Flink 看起来更完整,为什么还要创建 CEP 库?
使用 CEP 制作的模式比使用 Flink 标准 DataStream 操作符制作的模式更有效(更高的吞吐量/其他指标)?(如果可能,提供一些关于此的文章/论文/文档的链接)
【问题讨论】:
标签: scala apache-flink flink-cep