【问题标题】:flink cep sql Event Not triggeringflink cep sql 事件未触发
【发布时间】:2021-10-25 16:30:10
【问题描述】:

我在 Flink SQL 中使用 CEP 模式,它按预期工作,连接到 Kafka 代理。但是当我连接到基于集群的云 kafka 设置时,Flink CEP 没有触发。这是我的sql:

create table agent_action_detail 
(
    agent_id String, 
    room_id String, 
    create_time Bigint, 
    call_type String, 
    application_id String, 
    connect_time Bigint, 
    row_time TIMESTAMP_LTZ(3), WATERMARK for row_time as row_time  - INTERVAL '1' MINUTE) 
with ('connector'='kafka', 'topic'='agent-action-detail', ...)

然后我以json格式发送消息,例如

{"agent_id":"agent_221","room_id":"room1","create_time":1635206828877,"call_type":"inbound","application_id":"app1","connect_time":1635206501735,"row_time":"2021-10-25 16:07:09.019Z"}

在 flink web ui 中,水印可以正常工作 flink web ui

我运行我的 cep sql:

select * from agent_action_detail
 match_recognize(
    partition by agent_id 
    order by row_time 
    measures 
        last(BF.create_time) as create_time, 
        first(AF.connect_time) as connect_time 
    one row per match AFTER MATCH SKIP PAST LAST ROW 
    pattern (BF+ AF) define BF as BF.connect_time > 0 ,AF as AF.connect_time > 0
 )

每条 kafka 消息,connect_time > 0,但 flink 未触发。 有人可以帮忙解决这个问题吗,在此先感谢!


select * from agent_action_detail match_recognize( partition by agent_id order by row_time  measures AF.connect_time as connect_time one row per match pattern (BF AF) WITHIN INTERVAL '1' second define BF as (last(BF.connect_time, 1) < 1), AF as AF.connect_time >= 100)

这是另一个 cep sql 仍然无法正常工作。 而agent_action_detail表被另一个flink sql插入为

insert into agent_action_detail select data.agent_id, data.room_id, data.create_time, data.call_type, data.application_id, data.connect_time, now() from source_table where type = 'xxx'

【问题讨论】:

  • 我尝试了许多其他模式应该触发,但没有触发

标签: apache-flink flink-sql flink-cep


【解决方案1】:

有几件事会导致模式匹配不产生结果:

  • 输入实际上不包含模式
  • 水印处理不正确
  • 这种模式在某种程度上是病态的

这个特定的模式在没有退出条件的情况下循环。这种模式不允许模式匹配引擎的内部状态被清除,这会导致问题。

如果你直接使用 Flink CEP,我会告诉你 尝试添加until(condition)within(time) 来限制可能匹配的数量。

使用MATCH_RECOGNIZE,看看您是否可以在模式中添加一个独特的终止元素。


更新:由于您在修改模式后仍然没有得到任何结果,您应该确定水印是否是问题的根源。 CEP 依赖于按时间对输入流进行排序,这取决于水印——但前提是您使用的是事件时间。

最简单的测试方法是切换到使用处理时间:

create table agent_action_detail 
(
    agent_id String, 
    ...
    row_time AS PROCTIME()
)
with (...)

如果可行,那么时间戳或水印就是问题所在。例如,如果所有事件都迟到了,您将得不到任何结果。就您而言,我想知道 row_time 列中有任何数据。


如果这不能揭示问题,请分享一个可重现的最小示例,包括观察问题所需的数据。

【讨论】:

  • select * from ekyc_dashboard_agent_action_detail match_recognize( 按 agent_id 分区 order by row_time 将 AF.connect_time 测量为 connect_time 每个匹配模式一行 (BF AF) WITHIN INTERVAL '1' second define BF as (last(BF.connect_time , 1) = 100) 我可以帮忙看看这个sql吗?
  • 你想用不同于你的模式循环部分的东西来结束模式。这里的问题是 AF 和 BF 的定义完全相同。
  • 我绝对同意 BF as BF.connect_time > 0 ,AF as AF.connect_time > 0 是一样的。但是当我将我的 BF AF 更新为“BF as (last(BF.connect_time, 1) = 100)”时它仍然不起作用。我坐了整整两天,一无所知。 > <...>
  • 大卫,我使用 proctime() 并且我工作!非常感谢。然后我会继续尝试使用事件时间。现在我很开心~
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2016-12-10
  • 1970-01-01
  • 2019-06-10
  • 2023-03-27
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多