【发布时间】: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