You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

连接云托管Kafka集群时Flink SQL CEP事件未触发问题求助

使用Flink SQL的CEP Pattern对接单节点Kafka broker时运行符合预期,但切换到集群模式的云托管Kafka部署后,Flink CEP始终无法触发,相关代码及信息如下:

1. 建表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显示Watermark生成正常。

2. 触发测试用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,但规则始终未触发。

补充测试信息

额外测试了如下CEP SQL,同样无法触发:

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)

另外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'

内容的提问来源于stack exchange,提问作者slo

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.28 15:06:04