什么是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"}
第一条和第二条都没输出
第三条开始输出

更多推荐