Flink CEP SQL订单事件匹配及输出次数限制问题咨询
问题描述
业务场景
- 存在两个Kafka输入主题,Schema为:
eventName, ingestion_time(将用作Watermark), orderType, orderCountry - 主题数据示例:
- 第一个主题:
{"eventName": "orderCreated", "userId":123, "ingestionTime": "1665042169543", "orderType":"ecommerce","orderCountry": "UK"} - 第二个主题:
{"eventName": "orderSucess", "userId":123, "ingestionTime": "1665042189543", "orderType":"ecommerce","orderCountry": "USA"}
- 第一个主题:
需求
获取所有在5分钟窗口内触发orderCreated事件但未触发orderSucess事件的userId(按orderType、orderCountry分组),且每个用户针对同一orderType和orderCountry最多输出2次(即10分钟后移除该用户状态,不再输出)。
遇到的问题
- 不清楚Flink CEP SQL中
A not followed B的正确写法; - 不知道如何实现每个用户同一
orderType和orderCountry最多输出2次的限制,即连续2个5分钟窗口未触发第二个事件时移除状态。
当前尝试的SQL代码
SELECT * FROM union_event_table MATCH_RECOGNIZE( PARTITION BY orderType,orderCountry ORDER BY ingestion_time MEASURES A.userId as userId A.orderType as orderType A.orderCountry AS orderCountry ONE ROW PER MATCH PATTERN (A not followed B) WITHIN INTERVAL '5' MINUTES DEFINE A As A.eventName = 'orderCreated' B AS B.eventName = 'orderSucess' )
内容的提问来源于stack exchange,提问作者user9068199
相关产品推荐
相关产品推荐

