flink-cep详细解读
什么是Flink CEP?
CEP = Complex Event Processing 复杂事件处理。专门用来识别数据流中一组有顺序、有时间限制的事件组合,自定义匹配规则,抓取异常、风险、特定业务行为。
1.明确需求:
需求:监控登录,同一个人 10 分钟里连续 3 次输错密码(fail 登录),就弹出警告
2.代码具体操作
1)创建 Flink 运行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
2)造一堆登录测试数据
DataStreamSource<LoginEvent> ds1 = env.fromElements(一堆LoginEvent对象);
LoginEvent 里面存三样东西:用户 ID、登录结果(fail/success)、登录时间。
这里代替真实 Kafka 数据源,方便本地测试。
3)给每条数据打上真实时间 + 水位线(必须有)
SingleOutputStreamOperator<LoginEvent> ds2 = ds1.assignTimestampsAndWatermarks(...) extractTimestamp:把字符串时间2023-07-18 10:10:20转成数字毫秒时间戳
Watermark 水位线:告诉 Flink 这条数据真实发生时间,后面判断 10 分钟时长全靠这个
Duration.ZERO:不允许数据迟到,数据来了就认它的时间
4)keyBy:把每个人的数据分开算(超级关键)
KeyedStream<LoginEvent, String> keyedStream = ds2.keyBy(事件 -> 取用户id);
如果不分开:用户 1 的失败会算到用户 2 头上,彻底乱套
keyBy 之后:用户 1 一套独立计数器、用户 2 一套独立计数器,互不干扰。
5)Pattern:写死匹配规则(CEP 核心模板)
Pattern.begin("first", 策略)
.where(登录状态是fail)
.times(3)
.consecutive()
.within(10分钟)
begin("first", skipPastLastEvent())开始匹配,匹配成功后,从最后一条匹配数据的下一条重新计数,不会重复告警
.where(只看登录失败的数据)成功的数据直接跳过,不算次数
.times(3) 需要凑够 3 次失败
.consecutive() 必须严格连续中间一旦出现一次 success 成功,之前失败次数全部清零,从头重新数
.within(Time.minutes(10))这 3 次连续失败,整体时间跨度不能超过 10 分钟,超时作废
6)CEP.pattern:把数据流套进规则里
PatternStream<LoginEvent> resultStream = CEP.pattern(keyedStream, pattern);
让 Flink 用上面写好的规则,实时一条一条核对每个人的登录记录,
只要满足「连续 3 次失败、10 分钟内」,就把这 3 条失败数据存到 resultStream 里。
7)select:拿到匹配成功的数据,整理输出
resultStream.select(new PatternSelectFunction(){ public String select(Map<String, List<LoginEvent>> map){ //map里存着匹配成功的3条失败记录 //取出用户id返回打印 } }) map 结构:key 是我们 pattern 写的名字first, value 是装 3 条失败事件的集合代码里用 set 提取用户 id,最后输出用户编号
8)打印 + 启动程序
select.print(); env.execute(); print 控制台打印警告;execute 才真正启动整个任务跑流程。
3.疑问
这个 CEP 的整体流程就是:先创建 env 运行环境,接着接入业务数据,如果用到事件时间就需要配置水印处理乱序数据,之后该做什么呢?我们的需求是识别同一个用户 10 分钟内连续 3 次登录失败,匹配成功就输出告警的用户 ID,所以要根据用户 ID 做 keyBy 分组。接下来编写核心业务逻辑,借助 CEP 实现连续 3 次登录失败的规则识别,这里需要定义 pattern,我想问问 pattern 到底是什么?
我们会筛选事件的登录状态,只有状态为 fail 才会被匹配,通过
.times(3).consecutive().within(Time.minutes(10))分别用来限定匹配次数、要求事件严格连续,同时约束所有匹配事件要在 10 分钟的时间范围内。规则定义完成后,把分组后的数据流和我们写好的 Pattern 规则进行绑定,绑定结束后需要提取我们需要的业务数据。这里的 select 操作作用是什么?
解答:
理清楚整个流程:
创建执行环境 env
接入登录数据流(测试 / 真实 Kafka)
分配事件时间 + 水印(CEP 要算 10 分钟时间范围,必须有水印支撑 EventTime 时间判断)
keyBy(用户ID):每个用户单独一套匹配规则、单独计数,A 用户失败不影响 B 用户定义 Pattern:把你的业务规则(连续 3 次失败、10 分钟内)写成模板规则
CEP.pattern(分组流, pattern):数据流套上规则,自动筛选出符合条件的数据组
select():把匹配成功的一组原始事件,加工成你想要的告警文本 / 结果打印、启动任务
Pattern 到底是什么?
Pattern 是 Flink CEP 用来识别多事件序列行为的判定标准模板;CEP 本身不懂你的业务监控逻辑,必须依靠 Pattern 定义好 “什么样连续的事件组合属于异常风险”,没有 Pattern,CEP 无法完成序列匹配工作。
Pattern 是 Flink 给你封装好的链式语法骨架,固定一套调用结构,你往骨架里填自己业务条件就行。
基础通用骨架(万能模板):
Pattern<实体类,实体类> pattern = Pattern
// 1.匹配序列起点,配置匹配成功后的跳转策略
.begin("阶段标识名", 跳过策略)
// 2.当前事件过滤条件(单条事件筛选)
.where(自定义条件)
// 3.重复次数、是否严格连续
.times(次数).consecutive()/allowNonConsecutive()
// 4.整套序列最大时间有效期
.within(时间长度);
Pattern<LoginEvent, LoginEvent> pattern =
Pattern.<LoginEvent>begin("first", AfterMatchSkipStrategy.skipPastLastEvent())
.where(new SimpleCondition<LoginEvent>() {
@Override
public boolean filter(LoginEvent value) throws Exception {
return value.getStatus().equals("fail");
}
}).times(3).consecutive().within(Time.minutes(10));
为什么用select()?
PatternStream<LoginEvent> resultStream = CEP.pattern (keyedStream, pattern);
这一步搭建匹配通道,当数据匹配规则成功后会自动把匹配到的事件封装成
Map<String, List < 实体 >> 包裹起来,select()是一个转换函数接口,专门接收这个打包好的 map,将里面的数据改成自己想要的输出格式,可以自定义输出内容,它是 PatternStream 主流标准搭配,也有 flatSelect 等其他替换方案
4.完整的代码
package com.qyh.day08;
import lombok.SneakyThrows;
import org.apache.commons.lang3.time.DateUtils;
import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.cep.CEP;
import org.apache.flink.cep.PatternSelectFunction;
import org.apache.flink.cep.PatternStream;
import org.apache.flink.cep.nfa.aftermatch.AfterMatchSkipStrategy;
import org.apache.flink.cep.pattern.Pattern;
import org.apache.flink.cep.pattern.conditions.SimpleCondition;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.KeyedStream;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.time.Time;
import java.time.Duration;
import java.util.*;
public class Demo01 {
public static void main(String[] args) throws Exception {
//1. env-准备环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
DataStreamSource<LoginEvent> ds1 = env.fromElements(
new LoginEvent("1", "fail", "2023-07-18 10:10:20"),
new LoginEvent("1", "success", "2023-07-18 10:10:21"),
new LoginEvent("1", "fail", "2023-07-18 10:10:22"),
new LoginEvent("1", "fail", "2023-07-18 10:13:25"),
new LoginEvent("1", "fail", "2023-07-18 10:18:30"),
new LoginEvent("1", "fail", "2023-07-18 10:18:30"),
new LoginEvent("2", "fail", "2023-07-18 10:10:21"),
new LoginEvent("2", "fail", "2023-07-18 10:10:22"),
new LoginEvent("2", "success", "2023-07-18 10:10:23"),
new LoginEvent("2", "fail", "2023-07-18 10:10:24")
);
SingleOutputStreamOperator<LoginEvent> ds2 = ds1.assignTimestampsAndWatermarks(
WatermarkStrategy.<LoginEvent>forBoundedOutOfOrderness(Duration.ZERO)
.withTimestampAssigner(new SerializableTimestampAssigner<LoginEvent>() {
@SneakyThrows
@Override
public long extractTimestamp(LoginEvent element, long recordTimestamp) {
Date date = DateUtils.parseDate(element.getLoginTime(), "yyyy-MM-dd HH:mm:ss");
// 返回值是毫秒
return date.getTime();
}
})
);
KeyedStream<LoginEvent, String> keyedStream = ds2.keyBy(new KeySelector<LoginEvent, String>() {
@Override
public String getKey(LoginEvent value) throws Exception {
return value.getId();
}
});
Pattern<LoginEvent, LoginEvent> pattern = Pattern.<LoginEvent>begin("first", AfterMatchSkipStrategy.skipPastLastEvent())
.where(new SimpleCondition<LoginEvent>() {
@Override
public boolean filter(LoginEvent value) throws Exception {
return value.getStatus().equals("fail");
}
}).times(3).consecutive().within(Time.minutes(10));
PatternStream<LoginEvent> resultStream = CEP.pattern(keyedStream, pattern);
SingleOutputStreamOperator<String> select = resultStream.select(new PatternSelectFunction<LoginEvent, String>() {
@Override
public String select(Map<String, List<LoginEvent>> map) throws Exception {
Collection<List<LoginEvent>> values = map.values();
HashSet<String> set = new HashSet<>();
for (List<LoginEvent> list : values) {
for (LoginEvent loginEvent : list) {
set.add(loginEvent.getId());
}
}
return set.toString();
}
});
select.print("持续输错密码的账户:");
env.execute();
}
}
package com.qyh.day08;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@NoArgsConstructor
@AllArgsConstructor
public class LoginEvent {
private String id;
private String status;
private String loginTime;
}
运行结果:
5.另一个案例
1)需求
需求:同一个用户,先出现 create 事件 → 10 分钟内再出现 pay 事件,才算匹配成功;还要处理 10 分钟到了只收到 create、没收到 pay 的超时残缺数据
package com.qyh.day08;
import com.alibaba.fastjson.JSON;
import lombok.SneakyThrows;
import org.apache.commons.lang3.time.DateUtils;
import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.cep.CEP;
import org.apache.flink.cep.PatternSelectFunction;
import org.apache.flink.cep.PatternStream;
import org.apache.flink.cep.nfa.aftermatch.AfterMatchSkipStrategy;
import org.apache.flink.cep.pattern.Pattern;
import org.apache.flink.cep.pattern.conditions.SimpleCondition;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.KeyedStream;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.time.Time;
import java.time.Duration;
import java.util.*;
public class Demo02 {
public static void main(String[] args) throws Exception {
//案例二:两个事件,并且有超时情况的处理
// 需求:统计出数据中,第一个是 create,第二个是 pay ,并且这两个事件发生在 10 分钟之内
//1. env-准备环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
.setBootstrapServers("hadoop11:9092")
.setTopics("first1")
.setStartingOffsets(OffsetsInitializer.latest())
.setValueOnlyDeserializer(new SimpleStringSchema()).build();
DataStreamSource<String> ds1 = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "kafkaSource");
// JSON 字符串 转换成 PayEvent 实体类对象
SingleOutputStreamOperator<PayEvent> ds2 = ds1.map(new MapFunction<String, PayEvent>() {
@Override
public PayEvent map(String s) throws Exception {
return JSON.parseObject(s, PayEvent.class);
}
});
SingleOutputStreamOperator<PayEvent> ds3 = ds2.assignTimestampsAndWatermarks(
WatermarkStrategy.<PayEvent>forBoundedOutOfOrderness(Duration.ZERO)
.withTimestampAssigner(new SerializableTimestampAssigner<PayEvent>() {
@SneakyThrows
@Override
public long extractTimestamp(PayEvent element, long recordTimestamp) {
Date date = DateUtils.parseDate(element.getTs(), "yyyy-MM-dd HH:mm:ss");
// 返回值是毫秒
return date.getTime();
}
})
);
KeyedStream<PayEvent, Integer> keyedStream = ds3.keyBy(new KeySelector<PayEvent, Integer>() {
@Override
public Integer getKey(PayEvent value) throws Exception {
return value.getUserId();
}
});
Pattern<PayEvent, PayEvent> pattern = Pattern.<PayEvent>begin("create_step", AfterMatchSkipStrategy.skipPastLastEvent())
//第一步:匹配创建订单
.where(new SimpleCondition<PayEvent>() {
@Override
public boolean filter(PayEvent value) throws Exception {
return "create".equals(value.getType());
}
})
//严格下一个事件是支付
.next("pay_step")
.where(new SimpleCondition<PayEvent>() {
@Override
public boolean filter(PayEvent value) throws Exception {
return "pay".equals(value.getType());
}
})
//整体必须10分钟内完成
.within(Time.minutes(10));
PatternStream<PayEvent> resultStream = CEP.pattern(keyedStream, pattern);
resultStream.select(new PatternSelectFunction<PayEvent, Set<Integer>>() {
@Override
public Set<Integer> select(Map<String, List<PayEvent>> map) throws Exception {
HashSet<Integer> set = new HashSet<>();
Collection<List<PayEvent>> values = map.values();
for (List<PayEvent> value : values) {
for (PayEvent payEvent : value) {
set.add(payEvent.getUserId());
}
}
return set;
}
}).print();
env.execute();
}
}
2)需要的环境:
zookeeper Kafka KafkaUI界面
zk.sh satrt
kfc.sh start
kafkaUI.sh start UI界面:hadoop11:8889
3)测试
在first1中发送消息:
{"userId":1,"type":"create","ts":"2023-07-18 10:10:00"}
{"userId":1,"type":"pay","ts":"2023-07-18 10:15:00"}
{"userId":1,"type":"pay","ts":"2023-07-18 10:15:03"}
第一条和第二条都没输出
第三条开始输出
更多推荐


所有评论(0)