【问题标题】:Flink CEP not detecting last RecordFlink CEP 未检测到最后一条记录
【发布时间】:2020-02-25 20:43:06
【问题描述】:

我的代码有助于确定 Flink CEP 中记录的数量是否超过 25。因此,当我使用进程时间时,它匹配所有模式,但当我使用事件时间时,它与最后一条记录不匹配。

{"trasanction_id":196,"customer_id":27,"datetime":"1576499008876","amount":6094,"state":"SUCCESS"}
{"trasanction_id":197,"customer_id":27,"datetime":"1576499017565","amount":547,"state":"SUCCESS"}
{"trasanction_id":198,"customer_id":27,"datetime":"1576499029116","amount":6824,"state":"SUCCESS"}
{"trasanction_id":196,"customer_id":27,"datetime":"1576499053211","amount":6094,"state":"SUCCESS"}
{"trasanction_id":197,"customer_id":28,"datetime":"1576499063867","amount":547,"state":"FAILED"}
{"trasanction_id":198,"customer_id":28,"datetime":"1576499073566","amount":6824,"state":"FAILED"}

以上是我的记录。我有兴趣在事件时间匹配每个数量大于 25 的事件。理想情况下,它应该检测所有记录(它在处理时间中进行),因为所有记录的数量都大于 25。截至目前,我正在使用 3 秒的有界无序时间提取技术来实现无序。

请帮助我理解这一点。提前致谢! :)

【问题讨论】:

    标签: apache-flink complex-event-processing flink-cep


    【解决方案1】:

    由于 CEP 匹配时间模式,因此在使用事件时间时间戳时,事件首先按时间戳排序。这种排序涉及缓冲每个事件,直到水印赶上该事件,以便为任何较早的事件首先到达提供时间。

    由于您的水印配置为落后于流的前沿(即迄今为止最大的时间戳)3 秒,因此流的水印永远不会达到最后一个事件的时间戳。这就是没有处理最后一个事件的原因。 Flink 正在等待查看是否有任何更早的事件将到达,并且不会放弃,直到水印指示流已通过最后一个事件的时间戳完成。

    【讨论】:

    • 谢谢大卫。我的印象是 BoundedOutOfOrderness 在处理时间上起作用....但是,删除 BoundedOutOforderness 是唯一的解决方案吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-12-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多