Apache Flink CEP如何实现连续模式下字段值不同的登录异常检测
Flink CEP 实现该攻击检测场景的方案
前置准备:定义登录事件实体
首先统一登录事件的字段结构,示例采用Java实现:
// 登录事件类,需实现序列化接口 public class LoginEvent implements Serializable { // 用户名 private String username; // 设备ID private String deviceId; // 登录状态:SUCCESS/FAIL private String status; // 登录事件时间戳(毫秒) private Long eventTime; // 省略构造方法、getter、setter }
步骤1:配置事件时间与水位线
由于需要10分钟的时间窗口约束,优先使用事件时间避免处理时间漂移带来的误判:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 设置事件时间语义 env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); // 接入登录事件流后配置水位线,允许5秒的事件乱序延迟 DataStream<LoginEvent> loginStream = env.addSource(yourLoginEventSource) .assignTimestampsAndWatermarks( WatermarkStrategy.<LoginEvent>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) -> event.getEventTime()) );
步骤2:按用户名分流
检测逻辑以单个用户为维度,所以先按用户名字段做keyBy分组:
KeyedStream<LoginEvent, String> keyedStream = loginStream.keyBy(LoginEvent::getUsername);
步骤3:核心CEP匹配模式定义
这个是整个逻辑的核心,严格匹配「同一设备连续10次登录失败→其他设备登录成功」的时序,且全链路在10分钟内:
Pattern<LoginEvent, ?> attackPattern = Pattern // 匹配第一段:连续10次登录失败 .<LoginEvent>begin("failBatch") .where(event -> "FAIL".equals(event.getStatus())) .times(10) // 要求10次失败是严格连续的,中间不能穿插其他状态的同用户事件 .consecutive() // 要求10次失败都来自同一设备 .where(new IterativeCondition<LoginEvent>() { @Override public boolean filter(LoginEvent currentEvent, Context<LoginEvent> context) throws Exception { Iterable<LoginEvent> preFails = context.getEventsForPattern("failBatch"); for (LoginEvent preFail : preFails) { if (!preFail.getDeviceId().equals(currentEvent.getDeviceId())) { return false; } } return true; } }) // 匹配第二段:紧接着的登录成功事件 .next("successEvent") .where(event -> "SUCCESS".equals(event.getStatus())) // 要求成功事件的设备和前面失败批次的设备不一致 .where(new IterativeCondition<LoginEvent>() { @Override public boolean filter(LoginEvent successEvent, Context<LoginEvent> context) throws Exception { LoginEvent firstFail = context.getEventsForPattern("failBatch").iterator().next(); return !successEvent.getDeviceId().equals(firstFail.getDeviceId()); } }) // 整个事件链的时间窗口约束为10分钟 .within(Time.minutes(10));
步骤4:匹配结果提取与告警输出
将匹配到的规则事件解析为告警信息输出即可:
PatternStream<LoginEvent> patternStream = CEP.pattern(keyedStream, attackPattern); patternStream.process(new PatternProcessFunction<LoginEvent, Alert>() { @Override public void processMatch(Map<String, List<LoginEvent>> match, Context context, Collector<Alert> collector) throws Exception { List<LoginEvent> failBatch = match.get("failBatch"); LoginEvent successEvent = match.get("successEvent").get(0); // 构造告警对象 Alert alert = new Alert(); alert.setUsername(failBatch.get(0).getUsername()); alert.setFailDevice(failBatch.get(0).getDeviceId()); alert.setSuccessDevice(successEvent.getDeviceId()); alert.setFirstFailTime(failBatch.get(0).getEventTime()); alert.setLastFailTime(failBatch.get(failBatch.size()-1).getEventTime()); alert.setSuccessTime(successEvent.getEventTime()); collector.collect(alert); } }).print("攻击告警"); env.execute("LoginAttackDetection");
注意事项
- 水位线的乱序延迟时长可以根据实际生产环境的事件到达延迟情况调整,避免漏匹配
- 登录状态的枚举值可以根据实际业务的日志格式调整,比如部分日志会用0/1表示失败/成功
- 若需要放宽连续失败的限制,去掉
consecutive()配置即可匹配10分钟内非连续的同设备10次失败场景
内容的提问来源于stack exchange,提问作者JDForLife
相关产品推荐
相关产品推荐

