第一章.基于埋点日志的网络流量统计

1.数据准备

UserBehavior.csv

2.指定时间范围内网站总浏览量(PV)的统计

需求分析:实现一个网站总浏览量的统计,可以使用滚动时间窗口来实现,实时统计每小时内的网站PV

  1. package com.atguigu.flink.day07;
  2. import com.atguigu.flink.day07.pojo.UserBehavior;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.api.common.functions.AggregateFunction;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.streaming.api.functions.windowing.WindowFunction;
  7. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  8. import org.apache.flink.streaming.api.windowing.time.Time;
  9. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  10. import org.apache.flink.util.Collector;
  11. import java.time.Duration;
  12. import java.util.Date;
  13. /**
  14. * 需求:指定时间范围内网站总浏览量(PV)的统计
  15. */
  16. public class $01_PVCount {
  17. public static void main(String[] args) throws Exception {
  18. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  19. env.readTextFile("input/UserBehavior.csv")
  20. .map(line -> {
  21. String[] data = line.split(",");
  22. return new UserBehavior(
  23. Long.parseLong(data[0]),
  24. Long.parseLong(data[1]),
  25. Integer.parseInt(data[2]),
  26. data[3], Long.parseLong(data[4] ) * 1000
  27. );
  28. })
  29. .assignTimestampsAndWatermarks(
  30. WatermarkStrategy
  31. .<UserBehavior>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  32. .withTimestampAssigner((ub,ts)-> ub.getTimestamp())
  33. )
  34. .filter(ub -> "pv".equals(ub.getBehavior()))
  35. .keyBy(ub -> ub.getBehavior())
  36. .window(TumblingEventTimeWindows.of(Time.minutes(30)))
  37. .aggregate(
  38. new AggregateFunction<UserBehavior, Long, Long>() {
  39. @Override
  40. public Long createAccumulator() {
  41. return 0L;
  42. }
  43. @Override
  44. public Long add(UserBehavior value, Long accumulator) {
  45. return accumulator + 1;
  46. }
  47. @Override
  48. public Long getResult(Long accumulator) {
  49. return accumulator;
  50. }
  51. @Override
  52. public Long merge(Long a, Long b) {
  53. return null;
  54. }
  55. },
  56. new WindowFunction<Long, String, String, TimeWindow>() {
  57. @Override
  58. public void apply(String key, TimeWindow window, Iterable<Long> input, Collector<String> out) throws Exception {
  59. Long pv = input.iterator().next();
  60. String msg = new Date(window.getStart()) + " " + new Date(window.getEnd()) + " " + pv;
  61. out.collect(msg);
  62. }
  63. }
  64. )
  65. .print();
  66. env.execute();
  67. }
  68. }

day07[Flink流处理高阶编程实战] - 图1

3.指定时间范围内网站独立访客数(UV)的统计

  1. package com.atguigu.flink.day07;
  2. import com.atguigu.flink.day07.pojo.UserBehavior;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.api.common.functions.AggregateFunction;
  5. import org.apache.flink.api.common.state.MapState;
  6. import org.apache.flink.api.common.state.MapStateDescriptor;
  7. import org.apache.flink.api.common.typeinfo.Types;
  8. import org.apache.flink.configuration.Configuration;
  9. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  10. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  11. import org.apache.flink.streaming.api.functions.windowing.WindowFunction;
  12. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  13. import org.apache.flink.streaming.api.windowing.time.Time;
  14. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  15. import org.apache.flink.util.Collector;
  16. import java.time.Duration;
  17. import java.util.ArrayList;
  18. import java.util.Arrays;
  19. import java.util.Date;
  20. import java.util.Iterator;
  21. /**
  22. * 需求:指定时间范围内网站独立访客数(UV)的统计
  23. */
  24. public class $02_UVCount {
  25. public static void main(String[] args) throws Exception {
  26. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  27. env.readTextFile("input/UserBehavior.csv")
  28. .map(line -> {
  29. String[] data = line.split(",");
  30. return new UserBehavior(
  31. Long.parseLong(data[0]),
  32. Long.parseLong(data[1]),
  33. Integer.parseInt(data[2]),
  34. data[3], Long.parseLong(data[4] ) * 1000
  35. );
  36. })
  37. .assignTimestampsAndWatermarks(
  38. WatermarkStrategy
  39. .<UserBehavior>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  40. .withTimestampAssigner((ub,ts)-> ub.getTimestamp())
  41. )
  42. .filter(ub -> "pv".equals(ub.getBehavior()))
  43. .keyBy(ub -> ub.getBehavior())
  44. .window(TumblingEventTimeWindows.of(Time.minutes(30)))
  45. .process(new ProcessWindowFunction<UserBehavior, String, String, TimeWindow>() {
  46. private MapState<Long, String> userIdState;
  47. @Override
  48. public void open(Configuration parameters) throws Exception {
  49. userIdState = getRuntimeContext().getMapState(new MapStateDescriptor<Long, String>(
  50. "userIdState",
  51. Types.LONG,
  52. Types.STRING
  53. ));
  54. }
  55. @Override
  56. public void process(String s, Context context, Iterable<UserBehavior> elements, Collector<String> out) throws Exception {
  57. userIdState.clear();
  58. for (UserBehavior element : elements) {
  59. //所有用户的id存入到状态中,自动去重
  60. userIdState.put(element.getUserId(),"a");
  61. }
  62. Iterator<Long> iterator = userIdState.keys().iterator();
  63. ArrayList<Long> arrayList = new ArrayList<>();
  64. while(iterator.hasNext()){
  65. arrayList.add(iterator.next());
  66. }
  67. out.collect(context.window() + " " + arrayList.size());
  68. }
  69. })
  70. .print();
  71. env.execute();
  72. }
  73. }

day07[Flink流处理高阶编程实战] - 图2

第二章.电商数据分析

  1. 电商平台的用户行为频繁且比较复杂,系统上线运行一段时间后,可以收集到大量的用户行为数据,进而利用大数据技术进行深入挖掘和分析,得到感兴趣的商业指标并增强对风险的控制
  2. 电商用户行为数据多样,整体可以分为用户行为数据(日志数据)和业务行为数据两大类
  3. 用户的行为习惯数据包括了用户的登录方式,上线的时间点及时长,点击和浏览页面,页面停留时间以及页面跳转等等,我们可以从中进行流量统计和热门商品的统计,也可以深入挖掘用户的特征,这些数据往往可以从web服务器日志中直接读取到
  4. 二业务行为数据就是用户在电商平台中针对每个业务(通常是某个具体商品)所做的具体操作,我们一般会在业务系统中相应的位置埋点,然后收集日志进行分析

1.数据准备

UserBehavior.csv

2.实时热门商品统计

需求分析: 每隔30分钟输出最近一小时内点击量最多的前N个商品

  • 最近一小时:窗口长度
  • 每隔5分钟:窗口滑动步长
  • 时间:使用event-time(事件时间)

这里准备一个javabean用来保存商品信息

  1. package com.atguigu.flink.day07.pojo;
  2. import lombok.AllArgsConstructor;
  3. import lombok.Data;
  4. import lombok.NoArgsConstructor;
  5. @Data
  6. @AllArgsConstructor
  7. @NoArgsConstructor
  8. public class HotItem {
  9. private Long itemId;
  10. private Long windowEndTime;
  11. private Long count;
  12. }

具体实现代码

  1. package com.atguigu.flink.day07;
  2. import com.atguigu.flink.day07.pojo.HotItem;
  3. import com.atguigu.flink.day07.pojo.UserBehavior;
  4. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  5. import org.apache.flink.api.common.functions.AggregateFunction;
  6. import org.apache.flink.api.common.state.ListState;
  7. import org.apache.flink.api.common.state.ListStateDescriptor;
  8. import org.apache.flink.api.common.state.MapState;
  9. import org.apache.flink.api.common.state.MapStateDescriptor;
  10. import org.apache.flink.api.common.typeinfo.Types;
  11. import org.apache.flink.configuration.Configuration;
  12. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  13. import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
  14. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  15. import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows;
  16. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  17. import org.apache.flink.streaming.api.windowing.time.Time;
  18. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  19. import org.apache.flink.util.Collector;
  20. import java.time.Duration;
  21. import java.util.ArrayList;
  22. import java.util.Iterator;
  23. /**
  24. * 需求:每隔30分钟输出最近一小时内点击量最多的前N个商品
  25. */
  26. public class $03_TopN {
  27. public static void main(String[] args) throws Exception {
  28. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  29. env.readTextFile("input/UserBehavior.csv")
  30. .map(line -> {
  31. String[] data = line.split(",");
  32. return new UserBehavior(
  33. Long.parseLong(data[0]),
  34. Long.parseLong(data[1]),
  35. Integer.parseInt(data[2]),
  36. data[3], Long.parseLong(data[4] ) * 1000
  37. );
  38. })
  39. .assignTimestampsAndWatermarks(
  40. WatermarkStrategy
  41. .<UserBehavior>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  42. .withTimestampAssigner((ub,ts)-> ub.getTimestamp())
  43. )
  44. .filter(ub -> "pv".equals(ub.getBehavior()))
  45. .keyBy(ub -> ub.getItemId())
  46. .window(SlidingEventTimeWindows.of(Time.hours(1),Time.minutes(30)))
  47. .aggregate(
  48. new AggregateFunction<UserBehavior, Long, Long>() {
  49. @Override
  50. public Long createAccumulator() {
  51. return 0L;
  52. }
  53. @Override
  54. public Long add(UserBehavior value, Long accumulator) {
  55. return accumulator + 1;
  56. }
  57. @Override
  58. public Long getResult(Long accumulator) {
  59. return accumulator;
  60. }
  61. @Override
  62. public Long merge(Long a, Long b) {
  63. return null;
  64. }
  65. },
  66. new ProcessWindowFunction<Long, HotItem, Long, TimeWindow>() {
  67. @Override
  68. public void process(Long aLong, Context context, Iterable<Long> elements, Collector<HotItem> out) throws Exception {
  69. out.collect(new HotItem(aLong,context.window().getEnd(),elements.iterator().next()));
  70. }
  71. }
  72. )
  73. .keyBy(hotItem -> hotItem.getWindowEndTime())
  74. .process(new KeyedProcessFunction<Long, HotItem, String>() {
  75. private ListState<HotItem> hotItemState;
  76. @Override
  77. public void open(Configuration parameters) throws Exception {
  78. hotItemState = getRuntimeContext().getListState(new ListStateDescriptor<HotItem>(
  79. "hotItemState",
  80. HotItem.class
  81. ));
  82. }
  83. @Override
  84. public void processElement(HotItem value, Context ctx, Collector<String> out) throws Exception {
  85. //每来一条数据,把数据存入到状态,等到定时器触发的时候,进行排序
  86. if(!hotItemState.get().iterator().hasNext()){
  87. //现在是来了第一条数据:注册定时器
  88. ctx.timerService().registerEventTimeTimer(value.getWindowEndTime() + 5000L);
  89. }
  90. hotItemState.add(value);
  91. }
  92. @Override
  93. public void onTimer(long timestamp, OnTimerContext ctx, Collector<String> out) throws Exception {
  94. Iterator<HotItem> iterator = hotItemState.get().iterator();
  95. ArrayList<HotItem> hotItems = new ArrayList<>();
  96. while(iterator.hasNext()){
  97. hotItems.add(iterator.next());
  98. }
  99. hotItems.sort((o1,o2)->o2.getCount().compareTo(o1.getCount()));
  100. StringBuilder stringBuilder = new StringBuilder();
  101. stringBuilder.append("-------------------");
  102. for (int i = 0,count = Math.min(3, hotItems.size());i<count;i++){
  103. stringBuilder.append(hotItems.get(i)).append("\t");
  104. }
  105. out.collect(stringBuilder.toString());
  106. }
  107. })
  108. .print();
  109. env.execute();
  110. }
  111. }

day07[Flink流处理高阶编程实战] - 图3

第三章.页面广告分析

1.数据准备

AdClickLog.csv

2.黑名单过滤

  1. 我们进行的点击量计算,同一用户的重复点击是会叠加计算的
  2. 在实际场景中,同一用户确实可能反复点开同一个广告,这也说明了用户对广告更大的兴趣,但是如果用户在一段时间频繁的点击广告,这显然不是一个正常行为,有刷点击量的嫌疑
  3. 所以我们可以对一段时间内(比如一天内)的用户点击行为进行约束,如果对同一个广告点击超过一定限额(比如100次)应该把该用户加入黑名单并报警,此后其点击行为不应该再统计

两个功能

  • 告警:使用侧输出流
  • 已经进入黑名单的用户的广告点击记录不再进行统计
  1. package com.atguigu.flink.day07;
  2. import com.atguigu.flink.day07.pojo.AdsClickLog;
  3. import com.atguigu.flink.day07.pojo.HotItem;
  4. import com.atguigu.flink.day07.pojo.UserBehavior;
  5. import org.apache.commons.lang.math.LongRange;
  6. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  7. import org.apache.flink.api.common.functions.AggregateFunction;
  8. import org.apache.flink.api.common.functions.ReduceFunction;
  9. import org.apache.flink.api.common.state.*;
  10. import org.apache.flink.api.common.typeinfo.Types;
  11. import org.apache.flink.configuration.Configuration;
  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.functions.KeyedProcessFunction;
  15. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  16. import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows;
  17. import org.apache.flink.streaming.api.windowing.time.Time;
  18. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  19. import org.apache.flink.util.Collector;
  20. import org.apache.flink.util.OutputTag;
  21. import java.text.SimpleDateFormat;
  22. import java.time.Duration;
  23. import java.util.ArrayList;
  24. import java.util.Iterator;
  25. /**
  26. * 需求:黑名单过滤(时间间隔为一天,广告点击次数为100次)
  27. */
  28. public class $04_BlackList {
  29. public static void main(String[] args) throws Exception {
  30. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  31. env.setParallelism(1);
  32. SingleOutputStreamOperator<String> normal = env.readTextFile("input/AdClickLog.csv")
  33. .map(line -> {
  34. String[] data = line.split(",");
  35. return new AdsClickLog(
  36. Long.parseLong(data[0]),
  37. Long.parseLong(data[1]),
  38. data[2],
  39. data[3],
  40. Long.parseLong(data[4])
  41. );
  42. })
  43. .assignTimestampsAndWatermarks(
  44. WatermarkStrategy
  45. .<AdsClickLog>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  46. .withTimestampAssigner((log, ts) -> log.getTimestamp())
  47. )
  48. .keyBy(log -> log.getUserId() + ":" + log.getAdsId())
  49. .process(new KeyedProcessFunction<String, AdsClickLog, String>() {
  50. private ValueState<Boolean> warnState;
  51. private ReducingState<Long> countState;
  52. //定义一个状态存储年月日,如果新的数据和状态一致则是同一天,否则就是第二天
  53. private ValueState<String> yesterdayState;
  54. @Override
  55. public void open(Configuration parameters) throws Exception {
  56. countState = getRuntimeContext().getReducingState(
  57. new ReducingStateDescriptor<Long>(
  58. "countState",
  59. new ReduceFunction<Long>() {
  60. @Override
  61. public Long reduce(Long value1, Long value2) throws Exception {
  62. return value1 + value2;
  63. }
  64. },
  65. Types.LONG
  66. )
  67. );
  68. warnState = getRuntimeContext().getState(new ValueStateDescriptor<Boolean>("warnState", Types.BOOLEAN));
  69. yesterdayState = getRuntimeContext().getState(new ValueStateDescriptor<String>("yesterdayState", Types.STRING));
  70. }
  71. @Override
  72. public void processElement(AdsClickLog value, Context ctx, Collector<String> out) throws Exception {
  73. String today = new SimpleDateFormat("yyyyMMdd").format(value.getTimestamp());
  74. String yesterday = yesterdayState.value();
  75. //如果到了第二天,清空状态点击量和黑名单重新计算
  76. if (!today.equals(yesterday)) {
  77. countState.clear();
  78. warnState.clear();
  79. yesterdayState.update(today);
  80. }
  81. if (warnState.value() == null) {
  82. //说明没有预警,需要对用户的操作进行统计
  83. countState.add(1L);
  84. }
  85. String msg = value.getUserId() + "对广告:" + value.getAdsId() + "的点击量:" + countState.get();
  86. if (countState.get() > 99) {
  87. if (warnState.value() == null) {
  88. ctx.output(new OutputTag<String>("black") {
  89. }, msg + ",超过了阈值100,加入黑名单");
  90. warnState.update(true);
  91. }
  92. } else {
  93. out.collect(msg);
  94. }
  95. }
  96. });
  97. normal.print("正常");
  98. normal.getSideOutput(new OutputTag<String>("black"){}).print("black");
  99. env.execute();
  100. }
  101. }

day07[Flink流处理高阶编程实战] - 图4

第四章.恶意登录监控

  1. 对于网站而言,用户登录并不是频繁的业务操作,如果一个用户短时间内频繁登录失败,就有可能是出现了程序的恶意攻击,比如密码暴力破解
  2. 因此,我们考虑,应该对用户的登录失败动作进行统计,具体来说,如果一个同一用户在2s之内连续两次登录失败,就认为存在恶意登录的风险,输出相关的信息进行报警提示,这是电商网站,也是几乎所有网站风控的基本一环

1.数据准备

LoginLog.csv

2.封装数据的JavaBean

  1. package com.atguigu.flink.day07.pojo;
  2. import lombok.AllArgsConstructor;
  3. import lombok.Data;
  4. import lombok.NoArgsConstructor;
  5. @Data
  6. @NoArgsConstructor
  7. @AllArgsConstructor
  8. public class LoginEvent {
  9. private Long userId;
  10. private String ip;
  11. private String eventType;
  12. private Long eventTime;
  13. }

3.具体实现代码

  1. package com.atguigu.flink.day07;
  2. import com.atguigu.flink.day07.pojo.LoginEvent;
  3. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  4. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  5. import org.apache.flink.streaming.api.windowing.windows.GlobalWindow;
  6. import org.apache.flink.util.Collector;
  7. import java.util.ArrayList;
  8. import java.util.Iterator;
  9. /**
  10. * 需求:恶意登录监控(同一用户在2s之内连续两次登录失败)
  11. */
  12. public class $05_Login {
  13. public static void main(String[] args) throws Exception {
  14. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  15. env.setParallelism(1);
  16. env.readTextFile("input/LoginLog.csv")
  17. .map(line -> {
  18. String[] data = line.split(",");
  19. return new LoginEvent(
  20. Long.parseLong(data[0]),
  21. data[1],
  22. data[2],
  23. Long.parseLong(data[3] )* 1000L
  24. );
  25. })
  26. .keyBy(log->log.getUserId())
  27. .countWindow(2,1)
  28. .process(new ProcessWindowFunction<LoginEvent, String, Long, GlobalWindow>() {
  29. @Override
  30. public void process(Long aLong, Context context, Iterable<LoginEvent> elements, Collector<String> out) throws Exception {
  31. Iterator<LoginEvent> iterator = elements.iterator();
  32. ArrayList<LoginEvent> loginEvents = new ArrayList<>();
  33. while(iterator.hasNext()){
  34. loginEvents.add(iterator.next());
  35. }
  36. if(loginEvents.size() ==2){
  37. //如果长度为2才需要进行计算
  38. LoginEvent event1 = loginEvents.get(0);
  39. LoginEvent event2 = loginEvents.get(1);
  40. String eventType1 = event1.getEventType();
  41. String eventType2 = event2.getEventType();
  42. Long time1 = event1.getEventTime();
  43. Long time2 = event2.getEventTime();
  44. if("fail".equals(eventType1)
  45. && "fail".equals(eventType2)
  46. && Math.abs(time1 -time2)<=2000){
  47. out.collect(aLong + "在恶意登录,请注意....");
  48. }
  49. }
  50. }
  51. })
  52. .print();
  53. env.execute();
  54. }
  55. }

day07[Flink流处理高阶编程实战] - 图5

day07[Flink流处理高阶编程实战] - 图6

第五章.订单支付实时监控

  1. 在电商网站中,订单的支付作为直接与营销收入挂钩的一环,在业务流程中非常重要
  2. 对于订单而言,为了正确控制业务过程,也为了增加用户的支付意愿,网站一般会设置一个支付失效时间,超过一段时间不支付的订单就会被取消
  3. 另外,对于订单的支付,我们还应保证用户支付的正确性,这可以通过第三方支付平台的交易数据来做一个实时对账

1.数据准备

OrderLog.csv

2.封装数据的JavaBean

  1. package com.atguigu.flink.day07.pojo;
  2. import lombok.AllArgsConstructor;
  3. import lombok.Data;
  4. import lombok.NoArgsConstructor;
  5. @Data
  6. @AllArgsConstructor
  7. @NoArgsConstructor
  8. public class OrderEvent {
  9. private Long orderId;
  10. private String eventType;
  11. private String txId;
  12. private Long eventTime;
  13. }

3.具体代码实现

  1. package com.atguigu.flink.day07;
  2. import com.atguigu.flink.day07.pojo.OrderEvent;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.api.common.state.ValueState;
  5. import org.apache.flink.api.common.state.ValueStateDescriptor;
  6. import org.apache.flink.configuration.Configuration;
  7. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  8. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  9. import org.apache.flink.streaming.api.windowing.assigners.EventTimeSessionWindows;
  10. import org.apache.flink.streaming.api.windowing.time.Time;
  11. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  12. import org.apache.flink.util.Collector;
  13. import java.time.Duration;
  14. import java.util.ArrayList;
  15. import java.util.Iterator;
  16. /**
  17. * 需求:订单支付实时监控
  18. */
  19. public class $06_Order {
  20. public static void main(String[] args) throws Exception {
  21. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  22. env.setParallelism(1);
  23. env.readTextFile("input/OrderLog.csv")
  24. .map(line -> {
  25. String[] data = line.split(",");
  26. return new OrderEvent(
  27. Long.valueOf(data[0]),
  28. data[1],
  29. data[2],
  30. Long.parseLong(data[3]) * 1000
  31. );
  32. })
  33. .assignTimestampsAndWatermarks(
  34. WatermarkStrategy
  35. .<OrderEvent>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  36. .withTimestampAssigner((event, ts) -> event.getEventTime())
  37. )
  38. .keyBy(orderEvent -> orderEvent.getOrderId())
  39. .window(EventTimeSessionWindows.withGap(Time.minutes(45)))
  40. .process(new ProcessWindowFunction<OrderEvent, String, Long, TimeWindow>() {
  41. private ValueState<OrderEvent> creatState;
  42. @Override
  43. public void open(Configuration parameters) throws Exception {
  44. creatState = getRuntimeContext().getState(new ValueStateDescriptor<OrderEvent>(
  45. "creatState",
  46. OrderEvent.class
  47. ));
  48. }
  49. @Override
  50. public void process(Long key, Context context, Iterable<OrderEvent> elements, Collector<String> out) throws Exception {
  51. /**
  52. * 如果订单支付和创建都正常情况:则集合中有两个元素
  53. * 如果只有一个create,则表示没有key,或者pay的时间超过了45分钟
  54. */
  55. Iterator<OrderEvent> iterator = elements.iterator();
  56. ArrayList<OrderEvent> orderEvents = new ArrayList<>();
  57. while (iterator.hasNext()) {
  58. orderEvents.add(iterator.next());
  59. }
  60. if (orderEvents.size() == 2) {
  61. out.collect("订单:" + key + "正常");
  62. } else {
  63. //如果长度不是2,则一定是1
  64. OrderEvent event = orderEvents.get(0);
  65. if ("create".equals(event.getEventType())) {
  66. //这个session窗口只有create,没有pay,要么真的没有,要么超时
  67. creatState.update(event);
  68. } else {
  69. if (creatState.value() != null) {
  70. //表示create来过但是pay晚于create超过了45分钟,他们没有在一个窗口中
  71. out.collect("订单:" + key + "超时支付,请检查系统漏洞");
  72. } else {
  73. out.collect("订单:" + key + " 只有pay没有create, 请检测系统漏洞");
  74. }
  75. }
  76. }
  77. }
  78. })
  79. .print();
  80. env.execute();
  81. }
  82. }

day07[Flink流处理高阶编程实战] - 图7