第一章.Flink CEP介绍

1.什么是FlinkCEP

  1. FlinkCEP是在Flink实现的复杂事件处理库,它可以让你在无界流中检测出特定的数据,有机会掌握数据中重要的那部分
  2. 是一种基于动态环境中事件流的分析技术,事件在这里通常是有意义的状态变化,通过分析事件间的关系,利用过滤,联合聚合等技术,根据事件间的聚合关系和时序关系制定检测规则,持续的从事件流中查询出符合要求的事件序列,最终分析得到更复杂的复合事件
  3. 1.目标:从有序的简单事件流中发现一些高阶特征
  4. 2.输入:一个或多个由简单事件构成的事件流
  5. 3.处理:识别简单事件之间的内在联系,多个复合一定规则的简单事件构成复杂事件
  6. 4.输出:满足规则的复杂事件
  1. FlinkCEP应用场景
  2. 1.风险控制
  3. 2.策略营销
  4. 3.运维监控

2.Flink CEP开发

一.导入相关依赖

  1. <dependency>
  2. <groupId>org.apache.flink</groupId>
  3. <artifactId>flink-cep_${scala.binary.version}</artifactId>
  4. <version>${flink.version}</version>
  5. </dependency>

二.基本使用

  1. package com.atguigu.flink.day08;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.cep.CEP;
  5. import org.apache.flink.cep.PatternSelectFunction;
  6. import org.apache.flink.cep.PatternStream;
  7. import org.apache.flink.cep.nfa.aftermatch.AfterMatchSkipStrategy;
  8. import org.apache.flink.cep.pattern.Pattern;
  9. import org.apache.flink.cep.pattern.conditions.SimpleCondition;
  10. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  11. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  12. import java.time.Duration;
  13. import java.util.List;
  14. import java.util.Map;
  15. public class $01_CEPBaseUse {
  16. public static void main(String[] args) throws Exception {
  17. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  18. env.setParallelism(1);
  19. //1.先有数据流
  20. SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt")
  21. .map(line -> {
  22. String[] data = line.split(",");
  23. return new WaterSensor(
  24. data[0],
  25. Long.parseLong(data[1]),
  26. Integer.parseInt(data[2])
  27. );
  28. })
  29. .assignTimestampsAndWatermarks(
  30. WatermarkStrategy
  31. .<WaterSensor>forMonotonousTimestamps()
  32. .withTimestampAssigner((ws, ts) -> ws.getTs())
  33. );
  34. //2.定义模式
  35. Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start")
  36. .where(new SimpleCondition<WaterSensor>() {
  37. @Override
  38. public boolean filter(WaterSensor value) throws Exception {
  39. return "sensor_1".equals(value.getId());
  40. }
  41. });
  42. //3.把模式作用在数据流上
  43. PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);
  44. //4..获取匹配到的数据
  45. ps.select(new PatternSelectFunction<WaterSensor, String>() {
  46. @Override
  47. public String select(Map<String, List<WaterSensor>> map) throws Exception {
  48. return map.toString();
  49. }
  50. }).print();
  51. env.execute();
  52. }
  53. }

三.模式API

  1. 模式API可以让你定义想从输入流中抽取的复杂模式序列
  2. 模式:比如找拥有相同属性事件序列的模式
  3. 注意:1.每个模式必须有一个独一无二的名字,你可以在后面使用它来识别匹配到的事件
  4. 2.模式的名字不能包含字符":"
  5. 模式序列:每个复杂的模式序列包括多个简单的模式,也叫模式序列,你可以把模式序列看做是这样的模式构成的图,这些模式基于用户指定的条件从一个转换到另外一个
  6. 匹配:输入事件的一个序列,这些序列通过一系列有效的模式转换,能够访问到复杂模式图中的所有模式

1.单例模式

单例模式只接受一个事件,默认情况模式都是单例模式

上面的例子就是一个单例模式

2.组合模式

把多个单个模式组合在一起就是组合模式,组合模式由一个初始化模式(.begin(…))开头

  1. package com.atguigu.flink.day08;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.cep.CEP;
  5. import org.apache.flink.cep.PatternSelectFunction;
  6. import org.apache.flink.cep.PatternStream;
  7. import org.apache.flink.cep.pattern.Pattern;
  8. import org.apache.flink.cep.pattern.conditions.SimpleCondition;
  9. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  10. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  11. import java.util.List;
  12. import java.util.Map;
  13. public class $04_CEPBaseUseCombine {
  14. public static void main(String[] args) throws Exception {
  15. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  16. env.setParallelism(1);
  17. //1.先有数据流
  18. SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt")
  19. .map(line -> {
  20. String[] data = line.split(",");
  21. return new WaterSensor(
  22. data[0],
  23. Long.parseLong(data[1]),
  24. Integer.parseInt(data[2])
  25. );
  26. })
  27. .assignTimestampsAndWatermarks(
  28. WatermarkStrategy
  29. .<WaterSensor>forMonotonousTimestamps()
  30. .withTimestampAssigner((ws, ts) -> ws.getTs())
  31. );
  32. //2.定义模式
  33. Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start")
  34. .where(new SimpleCondition<WaterSensor>() {
  35. @Override
  36. public boolean filter(WaterSensor value) throws Exception {
  37. return "sensor_1".equals(value.getId());
  38. }
  39. })
  40. //.next("end") 严格连续
  41. //.followedBy("end") 松散连续
  42. .followedByAny("end") //不确定的松散连续
  43. .where(new SimpleCondition<WaterSensor>() {
  44. @Override
  45. public boolean filter(WaterSensor value) throws Exception {
  46. return "sensor_2".equals(value.getId());
  47. }
  48. });
  49. //3.把模式作用在数据流上
  50. PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);
  51. //4..获取匹配到的数据
  52. ps.select(new PatternSelectFunction<WaterSensor, String>() {
  53. @Override
  54. public String select(Map<String, List<WaterSensor>> map) throws Exception {
  55. return map.toString();
  56. }
  57. }).print();
  58. env.execute();
  59. }
  60. }

3.循环模式

循环模式可以接收多个事件

单例模式配合上量词就是循环模式(非常类似正则表达式)

  1. package com.atguigu.flink.day08;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.cep.CEP;
  5. import org.apache.flink.cep.PatternSelectFunction;
  6. import org.apache.flink.cep.PatternStream;
  7. import org.apache.flink.cep.pattern.Pattern;
  8. import org.apache.flink.cep.pattern.conditions.SimpleCondition;
  9. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  10. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  11. import java.util.List;
  12. import java.util.Map;
  13. public class $02_CEPBaseUseLoop {
  14. public static void main(String[] args) throws Exception {
  15. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  16. env.setParallelism(1);
  17. //1.先有数据流
  18. SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt")
  19. .map(line -> {
  20. String[] data = line.split(",");
  21. return new WaterSensor(
  22. data[0],
  23. Long.parseLong(data[1]),
  24. Integer.parseInt(data[2])
  25. );
  26. })
  27. .assignTimestampsAndWatermarks(
  28. WatermarkStrategy
  29. .<WaterSensor>forMonotonousTimestamps()
  30. .withTimestampAssigner((ws, ts) -> ws.getTs())
  31. );
  32. //2.定义模式
  33. Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start")
  34. .where(new SimpleCondition<WaterSensor>() {
  35. @Override
  36. public boolean filter(WaterSensor value) throws Exception {
  37. return "sensor_1".equals(value.getId());
  38. }
  39. })
  40. //.times(2) 出现两次;
  41. //.oneOrMore() //一次或多次
  42. //.times(2,4) //大于等于两次小于四次
  43. .timesOrMore(2)//2次或2次以上
  44. .until(new SimpleCondition<WaterSensor>() {
  45. @Override
  46. public boolean filter(WaterSensor value) throws Exception {
  47. return "sensor_2".equals(value.getId());
  48. }
  49. });
  50. //3.把模式作用在数据流上
  51. PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);
  52. //4..获取匹配到的数据
  53. ps.select(new PatternSelectFunction<WaterSensor, String>() {
  54. @Override
  55. public String select(Map<String, List<WaterSensor>> map) throws Exception {
  56. return map.toString();
  57. }
  58. }).print();
  59. env.execute();
  60. }
  61. }

4.循环模式的连续性

松散连续(默认)

忽略匹配的事件之间的不匹配的事件

  1. Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start")
  2. .where(new SimpleCondition<WaterSensor>() {
  3. @Override
  4. public boolean filter(WaterSensor value) throws Exception {
  5. return "sensor_1".equals(value.getId());
  6. }
  7. })
  8. .times(2); //出现两次

严格连续

期望所有匹配的事件严格的一个接一个出现,中间没有任何不匹配的事件

  1. Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start")
  2. .where(new SimpleCondition<WaterSensor>() {
  3. @Override
  4. public boolean filter(WaterSensor value) throws Exception {
  5. return "sensor_1".equals(value.getId());
  6. }
  7. })
  8. .times(2) //出现两次
  9. .consecutive();

非确定的松散连续

更近一步的松散连续,允许忽略掉一些匹配事件的附加匹配

  1. Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start")
  2. .where(new SimpleCondition<WaterSensor>() {
  3. @Override
  4. public boolean filter(WaterSensor value) throws Exception {
  5. return "sensor_1".equals(value.getId());
  6. }
  7. })
  8. .times(2) //出现两次
  9. .allowCombinations();

5.循环模式的贪婪性

在组合模式的情况下,对次数的处理尽可能的获取最多的那个个数,就是贪婪,当一个事件同时满足两个模式的时候起作用

  1. package com.atguigu.flink.day08;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.cep.CEP;
  5. import org.apache.flink.cep.PatternSelectFunction;
  6. import org.apache.flink.cep.PatternStream;
  7. import org.apache.flink.cep.pattern.Pattern;
  8. import org.apache.flink.cep.pattern.conditions.SimpleCondition;
  9. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  10. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  11. import java.util.List;
  12. import java.util.Map;
  13. public class $03_CEPBaseUseGreedy {
  14. public static void main(String[] args) throws Exception {
  15. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  16. env.setParallelism(1);
  17. //1.先有数据流
  18. SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt")
  19. .map(line -> {
  20. String[] data = line.split(",");
  21. return new WaterSensor(
  22. data[0],
  23. Long.parseLong(data[1]),
  24. Integer.parseInt(data[2])
  25. );
  26. })
  27. .assignTimestampsAndWatermarks(
  28. WatermarkStrategy
  29. .<WaterSensor>forMonotonousTimestamps()
  30. .withTimestampAssigner((ws, ts) -> ws.getTs())
  31. );
  32. //2.定义模式
  33. Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start")
  34. .where(new SimpleCondition<WaterSensor>() {
  35. @Override
  36. public boolean filter(WaterSensor value) throws Exception {
  37. return "sensor_1".equals(value.getId());
  38. }
  39. }).times(2,3).greedy()
  40. .next("end")
  41. .where(new SimpleCondition<WaterSensor>() {
  42. @Override
  43. public boolean filter(WaterSensor value) throws Exception {
  44. return value.getVc() == 30;
  45. }
  46. });
  47. //3.把模式作用在数据流上
  48. PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);
  49. //4..获取匹配到的数据
  50. ps.select(new PatternSelectFunction<WaterSensor, String>() {
  51. @Override
  52. public String select(Map<String, List<WaterSensor>> map) throws Exception {
  53. return map.toString();
  54. }
  55. }).print();
  56. env.execute();
  57. }
  58. }

6.模式的可选性

pattern.optional()方法让所有的模式变成可选的,不管是否是循环模式

  1. package com.atguigu.flink.day08;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.cep.CEP;
  5. import org.apache.flink.cep.PatternSelectFunction;
  6. import org.apache.flink.cep.PatternStream;
  7. import org.apache.flink.cep.pattern.Pattern;
  8. import org.apache.flink.cep.pattern.conditions.SimpleCondition;
  9. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  10. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  11. import java.util.List;
  12. import java.util.Map;
  13. public class $05_CEPBaseUseOptional {
  14. public static void main(String[] args) throws Exception {
  15. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  16. env.setParallelism(1);
  17. //1.先有数据流
  18. SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt")
  19. .map(line -> {
  20. String[] data = line.split(",");
  21. return new WaterSensor(
  22. data[0],
  23. Long.parseLong(data[1]),
  24. Integer.parseInt(data[2])
  25. );
  26. })
  27. .assignTimestampsAndWatermarks(
  28. WatermarkStrategy
  29. .<WaterSensor>forMonotonousTimestamps()
  30. .withTimestampAssigner((ws, ts) -> ws.getTs())
  31. );
  32. //2.定义模式
  33. Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start")
  34. .where(new SimpleCondition<WaterSensor>() {
  35. @Override
  36. public boolean filter(WaterSensor value) throws Exception {
  37. return "sensor_1".equals(value.getId());
  38. }
  39. }).times(2).optional() //出现0次或2次
  40. .next("end")
  41. .where(new SimpleCondition<WaterSensor>() {
  42. @Override
  43. public boolean filter(WaterSensor value) throws Exception {
  44. return "sensor_2".equals(value.getId());
  45. }
  46. });
  47. //3.把模式作用在数据流上
  48. PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);
  49. //4..获取匹配到的数据
  50. ps.select(new PatternSelectFunction<WaterSensor, String>() {
  51. @Override
  52. public String select(Map<String, List<WaterSensor>> map) throws Exception {
  53. return map.toString();
  54. }
  55. }).print();
  56. env.execute();
  57. }
  58. }

7.模式组

  1. 在前面的代码中次数只能用在某个模式上, 比如: .begin(...).where(...).next(...).where(...).times(2) 这里的次数只会用在next这个模式上, 而不会用在begin模式上.<br /> 如果需要用在多个模式上,可以使用模式组!
  1. package com.atguigu.flink.day08;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.cep.CEP;
  5. import org.apache.flink.cep.PatternSelectFunction;
  6. import org.apache.flink.cep.PatternStream;
  7. import org.apache.flink.cep.pattern.Pattern;
  8. import org.apache.flink.cep.pattern.conditions.SimpleCondition;
  9. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  10. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  11. import java.util.List;
  12. import java.util.Map;
  13. public class $06_CEPBaseUseGroup {
  14. public static void main(String[] args) throws Exception {
  15. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  16. env.setParallelism(1);
  17. //1.先有数据流
  18. SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt")
  19. .map(line -> {
  20. String[] data = line.split(",");
  21. return new WaterSensor(
  22. data[0],
  23. Long.parseLong(data[1]),
  24. Integer.parseInt(data[2])
  25. );
  26. })
  27. .assignTimestampsAndWatermarks(
  28. WatermarkStrategy
  29. .<WaterSensor>forMonotonousTimestamps()
  30. .withTimestampAssigner((ws, ts) -> ws.getTs())
  31. );
  32. //2.定义模式
  33. Pattern<WaterSensor, WaterSensor> pattern = Pattern.begin(Pattern.
  34. <WaterSensor>begin("start")
  35. .where(new SimpleCondition<WaterSensor>() {
  36. @Override
  37. public boolean filter(WaterSensor value) throws Exception {
  38. return "sensor_1".equals(value.getId());
  39. }
  40. })
  41. .next("next")
  42. .where(new SimpleCondition<WaterSensor>() {
  43. @Override
  44. public boolean filter(WaterSensor value) throws Exception {
  45. return "sensor_2".equals(value.getId());
  46. }
  47. }))
  48. .oneOrMore();
  49. //3.把模式作用在数据流上
  50. PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);
  51. //4..获取匹配到的数据
  52. ps.select(new PatternSelectFunction<WaterSensor, String>() {
  53. @Override
  54. public String select(Map<String, List<WaterSensor>> map) throws Exception {
  55. return map.toString();
  56. }
  57. }).print();
  58. env.execute();
  59. }
  60. }

8.超时数据

当一个模式上通过within加上窗口长度后,部分匹配的事件序列就可能因为超过窗口长度而被丢弃。

  1. package com.atguigu.flink.day08;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.cep.CEP;
  5. import org.apache.flink.cep.PatternSelectFunction;
  6. import org.apache.flink.cep.PatternStream;
  7. import org.apache.flink.cep.pattern.Pattern;
  8. import org.apache.flink.cep.pattern.conditions.SimpleCondition;
  9. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  10. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  11. import org.apache.flink.streaming.api.windowing.time.Time;
  12. import java.util.List;
  13. import java.util.Map;
  14. public class $07_CEPBaseUseWithIn {
  15. public static void main(String[] args) throws Exception {
  16. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  17. env.setParallelism(1);
  18. //1.先有数据流
  19. SingleOutputStreamOperator<WaterSensor> stream = env.readTextFile("input/sensor.txt")
  20. .map(line -> {
  21. String[] data = line.split(",");
  22. return new WaterSensor(
  23. data[0],
  24. Long.parseLong(data[1]),
  25. Integer.parseInt(data[2])
  26. );
  27. })
  28. .assignTimestampsAndWatermarks(
  29. WatermarkStrategy
  30. .<WaterSensor>forMonotonousTimestamps()
  31. .withTimestampAssigner((ws, ts) -> ws.getTs())
  32. );
  33. //2.定义模式
  34. Pattern<WaterSensor, WaterSensor> pattern = Pattern.<WaterSensor>begin("start")
  35. .where(new SimpleCondition<WaterSensor>() {
  36. @Override
  37. public boolean filter(WaterSensor value) throws Exception {
  38. return "sensor_1".equals(value.getId());
  39. }
  40. })
  41. .next("end")
  42. .where(new SimpleCondition<WaterSensor>() {
  43. @Override
  44. public boolean filter(WaterSensor value) throws Exception {
  45. return "sensor_2".equals(value.getId());
  46. }
  47. })
  48. .within(Time.seconds(2));
  49. //3.把模式作用在数据流上
  50. PatternStream<WaterSensor> ps = CEP.pattern(stream, pattern);
  51. //4..获取匹配到的数据
  52. ps.select(new PatternSelectFunction<WaterSensor, String>() {
  53. @Override
  54. public String select(Map<String, List<WaterSensor>> map) throws Exception {
  55. return map.toString();
  56. }
  57. }).print();
  58. env.execute();
  59. }
  60. }

9.匹配后跳过策略

day08[Flink CEP编程] - 图1

第二章.Flink CEP实战

1.恶意登录监控

  1. package com.atguigu.flink.day08;
  2. import com.atguigu.flink.day07.pojo.LoginEvent;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.cep.CEP;
  5. import org.apache.flink.cep.PatternSelectFunction;
  6. import org.apache.flink.cep.PatternStream;
  7. import org.apache.flink.cep.pattern.Pattern;
  8. import org.apache.flink.cep.pattern.conditions.SimpleCondition;
  9. import org.apache.flink.streaming.api.datastream.KeyedStream;
  10. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  11. import org.apache.flink.streaming.api.windowing.time.Time;
  12. import java.time.Duration;
  13. import java.util.List;
  14. import java.util.Map;
  15. /**
  16. * 需求:恶意登录监控(同一用户在2s之内连续两次登录失败)
  17. */
  18. public class $08_Login {
  19. public static void main(String[] args) throws Exception {
  20. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  21. env.setParallelism(1);
  22. KeyedStream<LoginEvent, Long> stream = env.readTextFile("input/LoginLog.csv")
  23. .map(line -> {
  24. String[] data = line.split(",");
  25. return new LoginEvent(
  26. Long.parseLong(data[0]),
  27. data[1],
  28. data[2],
  29. Long.parseLong(data[3]) * 1000L
  30. );
  31. })
  32. .assignTimestampsAndWatermarks(
  33. WatermarkStrategy.<LoginEvent>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  34. .withTimestampAssigner((log,ts)-> log.getEventTime())
  35. )
  36. .keyBy(log -> log.getUserId());
  37. //定义模式 连续两秒内登陆失败
  38. Pattern<LoginEvent, LoginEvent> pattern = Pattern
  39. .<LoginEvent>begin("fail_1")
  40. .where(new SimpleCondition<LoginEvent>() {
  41. @Override
  42. public boolean filter(LoginEvent value) throws Exception {
  43. return "fail".equals(value.getEventType());
  44. }
  45. })
  46. .next("fail_2")
  47. .where(new SimpleCondition<LoginEvent>() {
  48. @Override
  49. public boolean filter(LoginEvent value) throws Exception {
  50. return "fail".equals(value.getEventType());
  51. }
  52. })
  53. .within(Time.seconds(2001));
  54. //把模式作用在流上
  55. PatternStream<LoginEvent> ps = CEP.pattern(stream, pattern);
  56. //获取匹配到的数据
  57. ps.select(new PatternSelectFunction<LoginEvent, String>() {
  58. @Override
  59. public String select(Map<String, List<LoginEvent>> map) throws Exception {
  60. LoginEvent fail_1 = map.get("fail_1").get(0);
  61. return fail_1.getUserId() + "正在恶意登录,请注意";
  62. }
  63. })
  64. .print();
  65. env.execute();
  66. }
  67. }

2.订单支付实时监控

  1. package com.atguigu.flink.day08;
  2. import com.atguigu.flink.day07.pojo.OrderEvent;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.cep.CEP;
  5. import org.apache.flink.cep.PatternSelectFunction;
  6. import org.apache.flink.cep.PatternStream;
  7. import org.apache.flink.cep.PatternTimeoutFunction;
  8. import org.apache.flink.cep.nfa.aftermatch.AfterMatchSkipStrategy;
  9. import org.apache.flink.cep.pattern.Pattern;
  10. import org.apache.flink.cep.pattern.conditions.SimpleCondition;
  11. import org.apache.flink.streaming.api.datastream.KeyedStream;
  12. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  13. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  14. import org.apache.flink.streaming.api.windowing.time.Time;
  15. import org.apache.flink.util.OutputTag;
  16. import java.time.Duration;
  17. import java.util.List;
  18. import java.util.Map;
  19. /**
  20. * 需求:订单支付实时监控
  21. */
  22. public class $09_Order {
  23. public static void main(String[] args) throws Exception {
  24. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  25. env.setParallelism(1);
  26. KeyedStream<OrderEvent, Long> stream = env.readTextFile("input/OrderLog.csv")
  27. .map(line -> {
  28. String[] data = line.split(",");
  29. return new OrderEvent(
  30. Long.valueOf(data[0]),
  31. data[1],
  32. data[2],
  33. Long.parseLong(data[3]) * 1000
  34. );
  35. })
  36. .assignTimestampsAndWatermarks(
  37. WatermarkStrategy
  38. .<OrderEvent>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  39. .withTimestampAssigner((event, ts) -> event.getEventTime())
  40. )
  41. .keyBy(orderEvent -> orderEvent.getOrderId());
  42. Pattern<OrderEvent, OrderEvent> pattern = Pattern.<OrderEvent>begin("create", AfterMatchSkipStrategy.skipPastLastEvent())
  43. .where(new SimpleCondition<OrderEvent>() {
  44. @Override
  45. public boolean filter(OrderEvent value) throws Exception {
  46. return "create".equals(value.getEventType());
  47. }
  48. }).optional()
  49. .next("pay")
  50. .where(new SimpleCondition<OrderEvent>() {
  51. @Override
  52. public boolean filter(OrderEvent value) throws Exception {
  53. return "pay".equals(value.getEventType());
  54. }
  55. })
  56. .within(Time.minutes(45));
  57. PatternStream<OrderEvent> ps = CEP.pattern(stream, pattern);
  58. SingleOutputStreamOperator<String> normal = ps.select(
  59. new OutputTag<String>("late") {
  60. },
  61. new PatternTimeoutFunction<OrderEvent, String>() {
  62. @Override
  63. public String timeout(Map<String, List<OrderEvent>> map, long l) throws Exception {
  64. OrderEvent create = map.get("create").get(0);
  65. return "订单:" + create.getOrderId() + "只有create,没有pay,或者pay超时支付";
  66. }
  67. },
  68. new PatternSelectFunction<OrderEvent, String>() {
  69. @Override
  70. public String select(Map<String, List<OrderEvent>> map) throws Exception {
  71. OrderEvent pay = map.get("pay").get(0);
  72. if (!map.containsKey("create")) {
  73. return "订单:" + pay.getOrderId() + "只有pay,没有create,系统bug";
  74. } else {
  75. return "";
  76. }
  77. }
  78. }
  79. );
  80. normal.filter(s->s.length()>0)
  81. .union(normal.getSideOutput(new OutputTag<String>("late"){}))
  82. .print();
  83. env.execute();
  84. }
  85. }