【问题标题】:What is the difference between Flink Core and Flink CEP in respect of their capabilities?Flink Core 和 Flink CEP 在能力方面有什么区别?
【发布时间】: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

尽管这些策略无法生成所需的模式,但可以使用RichFunctionValueState 计数器轻松确定何时应发出新警报,将输入流转换为事件流。

因此,我希望对这些问题有所了解:

  • 如果 Flink 看起来更完整,为什么还要创建 CEP 库?

  • 使用 CEP 制作的模式比使用 Flink 标准 DataStream 操作符制作的模式更有效(更高的吞吐量/其他指标)?(如果可能,提供一些关于此的文章/论文/文档的链接)

【问题讨论】:

    标签: scala apache-flink flink-cep


    【解决方案1】:

    感谢您使用 Flink CEP。

    Flink CEP 是 Fl​​ink 之上的一个库。因此,它不会添加任何无法使用 vanilla Flink(ProcessFunctions 等)实现的功能。事实上,在底层,它被实现为一个特殊的运算符,它检查与特定模式匹配的元素,并且它的大部分功能甚至可能被实现为 ProcessFunction(周围有很多工具)。

    也就是说,Flink CEP 可能不会添加普通 Flink 无法实现的功能,但它增加了表达性,使得一些用例更易于实现。其他 API 也是如此,例如 Flink 中的 Windowing API,您可以使用 ProcessFunctions 来实现它(周围有很多工具)。

    现在谈到效率,答案是“视情况而定”。手工制作一个针对您的用例定制的特殊过程函数,并针对您的工作负载进行所有可能的优化,这可能比 FlinkCEP 更有效,因为后者是一个通用库。如果您有专业知识和时间,那么最佳解决方案始终是同时使用(CEP 和 vanilla Flink)实现 PoC,并为您的案例选择最有效的。

    【讨论】:

    • 感谢您的回答,这正是我的想法。我今天终于找到了this,这与我对 Flink CEP 的印象一致。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多