第一章.Flink CEP介绍
1.什么是FlinkCEP
FlinkCEP是在Flink实现的复杂事件处理库,它可以让你在无界流中检测出特定的数据,有机会掌握数据中重要的那部分是一种基于动态环境中事件流的分析技术,事件在这里通常是有意义的状态变化,通过分析事件间的关系,利用过滤,联合聚合等技术,根据事件间的聚合关系和时序关系制定检测规则,持续的从事件流中查询出符合要求的事件序列,最终分析得到更复杂的复合事件1.目标:从有序的简单事件流中发现一些高阶特征2.输入:一个或多个由简单事件构成的事件流3.处理:识别简单事件之间的内在联系,多个复合一定规则的简单事件构成复杂事件4.输出:满足规则的复杂事件
FlinkCEP应用场景1.风险控制2.策略营销3.运维监控
2.Flink CEP开发
一.导入相关依赖
<dependency><groupId>org.apache.flink</groupId><artifactId>flink-cep_${scala.binary.version}</artifactId><version>${flink.version}</version></dependency>
二.基本使用
package com.atguigu.flink.day08;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.WatermarkStrategy;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.SingleOutputStreamOperator;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import java.time.Duration;import java.util.List;import java.util.Map;public class $01_CEPBaseUse {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);//1.先有数据流SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt").map(line -> {String[] data = line.split(",");return new WaterSensor(data[0],Long.parseLong(data[1]),Integer.parseInt(data[2]));}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forMonotonousTimestamps().withTimestampAssigner((ws, ts) -> ws.getTs()));//2.定义模式Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_1".equals(value.getId());}});//3.把模式作用在数据流上PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);//4..获取匹配到的数据ps.select(new PatternSelectFunction<WaterSensor, String>() {@Overridepublic String select(Map<String, List<WaterSensor>> map) throws Exception {return map.toString();}}).print();env.execute();}}
三.模式API
模式API可以让你定义想从输入流中抽取的复杂模式序列模式:比如找拥有相同属性事件序列的模式注意:1.每个模式必须有一个独一无二的名字,你可以在后面使用它来识别匹配到的事件2.模式的名字不能包含字符":"模式序列:每个复杂的模式序列包括多个简单的模式,也叫模式序列,你可以把模式序列看做是这样的模式构成的图,这些模式基于用户指定的条件从一个转换到另外一个匹配:输入事件的一个序列,这些序列通过一系列有效的模式转换,能够访问到复杂模式图中的所有模式
1.单例模式
单例模式只接受一个事件,默认情况模式都是单例模式
上面的例子就是一个单例模式
2.组合模式
把多个单个模式组合在一起就是组合模式,组合模式由一个初始化模式(.begin(…))开头
package com.atguigu.flink.day08;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.cep.CEP;import org.apache.flink.cep.PatternSelectFunction;import org.apache.flink.cep.PatternStream;import org.apache.flink.cep.pattern.Pattern;import org.apache.flink.cep.pattern.conditions.SimpleCondition;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import java.util.List;import java.util.Map;public class $04_CEPBaseUseCombine {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);//1.先有数据流SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt").map(line -> {String[] data = line.split(",");return new WaterSensor(data[0],Long.parseLong(data[1]),Integer.parseInt(data[2]));}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forMonotonousTimestamps().withTimestampAssigner((ws, ts) -> ws.getTs()));//2.定义模式Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_1".equals(value.getId());}})//.next("end") 严格连续//.followedBy("end") 松散连续.followedByAny("end") //不确定的松散连续.where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_2".equals(value.getId());}});//3.把模式作用在数据流上PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);//4..获取匹配到的数据ps.select(new PatternSelectFunction<WaterSensor, String>() {@Overridepublic String select(Map<String, List<WaterSensor>> map) throws Exception {return map.toString();}}).print();env.execute();}}
3.循环模式
循环模式可以接收多个事件
单例模式配合上量词就是循环模式(非常类似正则表达式)
package com.atguigu.flink.day08;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.cep.CEP;import org.apache.flink.cep.PatternSelectFunction;import org.apache.flink.cep.PatternStream;import org.apache.flink.cep.pattern.Pattern;import org.apache.flink.cep.pattern.conditions.SimpleCondition;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import java.util.List;import java.util.Map;public class $02_CEPBaseUseLoop {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);//1.先有数据流SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt").map(line -> {String[] data = line.split(",");return new WaterSensor(data[0],Long.parseLong(data[1]),Integer.parseInt(data[2]));}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forMonotonousTimestamps().withTimestampAssigner((ws, ts) -> ws.getTs()));//2.定义模式Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_1".equals(value.getId());}})//.times(2) 出现两次;//.oneOrMore() //一次或多次//.times(2,4) //大于等于两次小于四次.timesOrMore(2)//2次或2次以上.until(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_2".equals(value.getId());}});//3.把模式作用在数据流上PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);//4..获取匹配到的数据ps.select(new PatternSelectFunction<WaterSensor, String>() {@Overridepublic String select(Map<String, List<WaterSensor>> map) throws Exception {return map.toString();}}).print();env.execute();}}
4.循环模式的连续性
松散连续(默认)
忽略匹配的事件之间的不匹配的事件
Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_1".equals(value.getId());}}).times(2); //出现两次
严格连续
期望所有匹配的事件严格的一个接一个出现,中间没有任何不匹配的事件
Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_1".equals(value.getId());}}).times(2) //出现两次.consecutive();
非确定的松散连续
更近一步的松散连续,允许忽略掉一些匹配事件的附加匹配
Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_1".equals(value.getId());}}).times(2) //出现两次.allowCombinations();
5.循环模式的贪婪性
在组合模式的情况下,对次数的处理尽可能的获取最多的那个个数,就是贪婪,当一个事件同时满足两个模式的时候起作用
package com.atguigu.flink.day08;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.cep.CEP;import org.apache.flink.cep.PatternSelectFunction;import org.apache.flink.cep.PatternStream;import org.apache.flink.cep.pattern.Pattern;import org.apache.flink.cep.pattern.conditions.SimpleCondition;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import java.util.List;import java.util.Map;public class $03_CEPBaseUseGreedy {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);//1.先有数据流SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt").map(line -> {String[] data = line.split(",");return new WaterSensor(data[0],Long.parseLong(data[1]),Integer.parseInt(data[2]));}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forMonotonousTimestamps().withTimestampAssigner((ws, ts) -> ws.getTs()));//2.定义模式Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_1".equals(value.getId());}}).times(2,3).greedy().next("end").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return value.getVc() == 30;}});//3.把模式作用在数据流上PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);//4..获取匹配到的数据ps.select(new PatternSelectFunction<WaterSensor, String>() {@Overridepublic String select(Map<String, List<WaterSensor>> map) throws Exception {return map.toString();}}).print();env.execute();}}
6.模式的可选性
pattern.optional()方法让所有的模式变成可选的,不管是否是循环模式
package com.atguigu.flink.day08;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.cep.CEP;import org.apache.flink.cep.PatternSelectFunction;import org.apache.flink.cep.PatternStream;import org.apache.flink.cep.pattern.Pattern;import org.apache.flink.cep.pattern.conditions.SimpleCondition;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import java.util.List;import java.util.Map;public class $05_CEPBaseUseOptional {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);//1.先有数据流SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt").map(line -> {String[] data = line.split(",");return new WaterSensor(data[0],Long.parseLong(data[1]),Integer.parseInt(data[2]));}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forMonotonousTimestamps().withTimestampAssigner((ws, ts) -> ws.getTs()));//2.定义模式Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_1".equals(value.getId());}}).times(2).optional() //出现0次或2次.next("end").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_2".equals(value.getId());}});//3.把模式作用在数据流上PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);//4..获取匹配到的数据ps.select(new PatternSelectFunction<WaterSensor, String>() {@Overridepublic String select(Map<String, List<WaterSensor>> map) throws Exception {return map.toString();}}).print();env.execute();}}
7.模式组
在前面的代码中次数只能用在某个模式上, 比如: .begin(...).where(...).next(...).where(...).times(2) 这里的次数只会用在next这个模式上, 而不会用在begin模式上.<br /> 如果需要用在多个模式上,可以使用模式组!
package com.atguigu.flink.day08;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.cep.CEP;import org.apache.flink.cep.PatternSelectFunction;import org.apache.flink.cep.PatternStream;import org.apache.flink.cep.pattern.Pattern;import org.apache.flink.cep.pattern.conditions.SimpleCondition;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import java.util.List;import java.util.Map;public class $06_CEPBaseUseGroup {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);//1.先有数据流SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt").map(line -> {String[] data = line.split(",");return new WaterSensor(data[0],Long.parseLong(data[1]),Integer.parseInt(data[2]));}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forMonotonousTimestamps().withTimestampAssigner((ws, ts) -> ws.getTs()));//2.定义模式Pattern<WaterSensor, WaterSensor> pattern = Pattern.begin(Pattern.<WaterSensor>begin("start").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_1".equals(value.getId());}}).next("next").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_2".equals(value.getId());}})).oneOrMore();//3.把模式作用在数据流上PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);//4..获取匹配到的数据ps.select(new PatternSelectFunction<WaterSensor, String>() {@Overridepublic String select(Map<String, List<WaterSensor>> map) throws Exception {return map.toString();}}).print();env.execute();}}
8.超时数据
当一个模式上通过within加上窗口长度后,部分匹配的事件序列就可能因为超过窗口长度而被丢弃。
package com.atguigu.flink.day08;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.cep.CEP;import org.apache.flink.cep.PatternSelectFunction;import org.apache.flink.cep.PatternStream;import org.apache.flink.cep.pattern.Pattern;import org.apache.flink.cep.pattern.conditions.SimpleCondition;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.util.List;import java.util.Map;public class $07_CEPBaseUseWithIn {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);//1.先有数据流SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt").map(line -> {String[] data = line.split(",");return new WaterSensor(data[0],Long.parseLong(data[1]),Integer.parseInt(data[2]));}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forMonotonousTimestamps().withTimestampAssigner((ws, ts) -> ws.getTs()));//2.定义模式Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_1".equals(value.getId());}}).next("end").where(new SimpleCondition<WaterSensor>() {@Overridepublic boolean filter(WaterSensor value) throws Exception {return "sensor_2".equals(value.getId());}}).within(Time.seconds(2));//3.把模式作用在数据流上PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);//4..获取匹配到的数据ps.select(new PatternSelectFunction<WaterSensor, String>() {@Overridepublic String select(Map<String, List<WaterSensor>> map) throws Exception {return map.toString();}}).print();env.execute();}}
9.匹配后跳过策略
![day08[Flink CEP编程] - 图1](/uploads/projects/liuye-6lcqc@ddtw8t/e8761dd4fc7972f86cebce000e38e563.png)
第二章.Flink CEP实战
1.恶意登录监控
package com.atguigu.flink.day08;import com.atguigu.flink.day07.pojo.LoginEvent;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.cep.CEP;import org.apache.flink.cep.PatternSelectFunction;import org.apache.flink.cep.PatternStream;import org.apache.flink.cep.pattern.Pattern;import org.apache.flink.cep.pattern.conditions.SimpleCondition;import org.apache.flink.streaming.api.datastream.KeyedStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.windowing.time.Time;import java.time.Duration;import java.util.List;import java.util.Map;/*** 需求:恶意登录监控(同一用户在2s之内连续两次登录失败)*/public class $08_Login {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);KeyedStream<LoginEvent, Long> stream = env.readTextFile("input/LoginLog.csv").map(line -> {String[] data = line.split(",");return new LoginEvent(Long.parseLong(data[0]),data[1],data[2],Long.parseLong(data[3]) * 1000L);}).assignTimestampsAndWatermarks(WatermarkStrategy.<LoginEvent>forBoundedOutOfOrderness(Duration.ofSeconds(3)).withTimestampAssigner((log,ts)-> log.getEventTime())).keyBy(log -> log.getUserId());//定义模式 连续两秒内登陆失败Pattern<LoginEvent, LoginEvent> pattern = Pattern.<LoginEvent>begin("fail_1").where(new SimpleCondition<LoginEvent>() {@Overridepublic boolean filter(LoginEvent value) throws Exception {return "fail".equals(value.getEventType());}}).next("fail_2").where(new SimpleCondition<LoginEvent>() {@Overridepublic boolean filter(LoginEvent value) throws Exception {return "fail".equals(value.getEventType());}}).within(Time.seconds(2001));//把模式作用在流上PatternStream<LoginEvent> ps = CEP.pattern(stream, pattern);//获取匹配到的数据ps.select(new PatternSelectFunction<LoginEvent, String>() {@Overridepublic String select(Map<String, List<LoginEvent>> map) throws Exception {LoginEvent fail_1 = map.get("fail_1").get(0);return fail_1.getUserId() + "正在恶意登录,请注意";}}).print();env.execute();}}
2.订单支付实时监控
package com.atguigu.flink.day08;import com.atguigu.flink.day07.pojo.OrderEvent;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.cep.CEP;import org.apache.flink.cep.PatternSelectFunction;import org.apache.flink.cep.PatternStream;import org.apache.flink.cep.PatternTimeoutFunction;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.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 org.apache.flink.util.OutputTag;import java.time.Duration;import java.util.List;import java.util.Map;/*** 需求:订单支付实时监控*/public class $09_Order {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);KeyedStream<OrderEvent, Long> stream = env.readTextFile("input/OrderLog.csv").map(line -> {String[] data = line.split(",");return new OrderEvent(Long.valueOf(data[0]),data[1],data[2],Long.parseLong(data[3]) * 1000);}).assignTimestampsAndWatermarks(WatermarkStrategy.<OrderEvent>forBoundedOutOfOrderness(Duration.ofSeconds(3)).withTimestampAssigner((event, ts) -> event.getEventTime())).keyBy(orderEvent -> orderEvent.getOrderId());Pattern<OrderEvent, OrderEvent> pattern = Pattern.<OrderEvent>begin("create", AfterMatchSkipStrategy.skipPastLastEvent()).where(new SimpleCondition<OrderEvent>() {@Overridepublic boolean filter(OrderEvent value) throws Exception {return "create".equals(value.getEventType());}}).optional().next("pay").where(new SimpleCondition<OrderEvent>() {@Overridepublic boolean filter(OrderEvent value) throws Exception {return "pay".equals(value.getEventType());}}).within(Time.minutes(45));PatternStream<OrderEvent> ps = CEP.pattern(stream, pattern);SingleOutputStreamOperator<String> normal = ps.select(new OutputTag<String>("late") {},new PatternTimeoutFunction<OrderEvent, String>() {@Overridepublic String timeout(Map<String, List<OrderEvent>> map, long l) throws Exception {OrderEvent create = map.get("create").get(0);return "订单:" + create.getOrderId() + "只有create,没有pay,或者pay超时支付";}},new PatternSelectFunction<OrderEvent, String>() {@Overridepublic String select(Map<String, List<OrderEvent>> map) throws Exception {OrderEvent pay = map.get("pay").get(0);if (!map.containsKey("create")) {return "订单:" + pay.getOrderId() + "只有pay,没有create,系统bug";} else {return "";}}});normal.filter(s->s.length()>0).union(normal.getSideOutput(new OutputTag<String>("late"){})).print();env.execute();}}
