第一章.基于埋点日志的网络流量统计
1.数据准备
2.指定时间范围内网站总浏览量(PV)的统计
需求分析:实现一个网站总浏览量的统计,可以使用滚动时间窗口来实现,实时统计每小时内的网站PV
package com.atguigu.flink.day07;import com.atguigu.flink.day07.pojo.UserBehavior;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.AggregateFunction;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.windowing.WindowFunction;import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import java.time.Duration;import java.util.Date;/*** 需求:指定时间范围内网站总浏览量(PV)的统计*/public class $01_PVCount {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.readTextFile("input/UserBehavior.csv").map(line -> {String[] data = line.split(",");return new UserBehavior(Long.parseLong(data[0]),Long.parseLong(data[1]),Integer.parseInt(data[2]),data[3], Long.parseLong(data[4] ) * 1000);}).assignTimestampsAndWatermarks(WatermarkStrategy.<UserBehavior>forBoundedOutOfOrderness(Duration.ofSeconds(3)).withTimestampAssigner((ub,ts)-> ub.getTimestamp())).filter(ub -> "pv".equals(ub.getBehavior())).keyBy(ub -> ub.getBehavior()).window(TumblingEventTimeWindows.of(Time.minutes(30))).aggregate(new AggregateFunction<UserBehavior, Long, Long>() {@Overridepublic Long createAccumulator() {return 0L;}@Overridepublic Long add(UserBehavior value, Long accumulator) {return accumulator + 1;}@Overridepublic Long getResult(Long accumulator) {return accumulator;}@Overridepublic Long merge(Long a, Long b) {return null;}},new WindowFunction<Long, String, String, TimeWindow>() {@Overridepublic void apply(String key, TimeWindow window, Iterable<Long> input, Collector<String> out) throws Exception {Long pv = input.iterator().next();String msg = new Date(window.getStart()) + " " + new Date(window.getEnd()) + " " + pv;out.collect(msg);}}).print();env.execute();}}
![day07[Flink流处理高阶编程实战] - 图1](/uploads/projects/liuye-6lcqc@ddtw8t/39761179bad369c116e3407491be862c.png)
3.指定时间范围内网站独立访客数(UV)的统计
package com.atguigu.flink.day07;import com.atguigu.flink.day07.pojo.UserBehavior;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.AggregateFunction;import org.apache.flink.api.common.state.MapState;import org.apache.flink.api.common.state.MapStateDescriptor;import org.apache.flink.api.common.typeinfo.Types;import org.apache.flink.configuration.Configuration;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.functions.windowing.WindowFunction;import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import java.time.Duration;import java.util.ArrayList;import java.util.Arrays;import java.util.Date;import java.util.Iterator;/*** 需求:指定时间范围内网站独立访客数(UV)的统计*/public class $02_UVCount {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.readTextFile("input/UserBehavior.csv").map(line -> {String[] data = line.split(",");return new UserBehavior(Long.parseLong(data[0]),Long.parseLong(data[1]),Integer.parseInt(data[2]),data[3], Long.parseLong(data[4] ) * 1000);}).assignTimestampsAndWatermarks(WatermarkStrategy.<UserBehavior>forBoundedOutOfOrderness(Duration.ofSeconds(3)).withTimestampAssigner((ub,ts)-> ub.getTimestamp())).filter(ub -> "pv".equals(ub.getBehavior())).keyBy(ub -> ub.getBehavior()).window(TumblingEventTimeWindows.of(Time.minutes(30))).process(new ProcessWindowFunction<UserBehavior, String, String, TimeWindow>() {private MapState<Long, String> userIdState;@Overridepublic void open(Configuration parameters) throws Exception {userIdState = getRuntimeContext().getMapState(new MapStateDescriptor<Long, String>("userIdState",Types.LONG,Types.STRING));}@Overridepublic void process(String s, Context context, Iterable<UserBehavior> elements, Collector<String> out) throws Exception {userIdState.clear();for (UserBehavior element : elements) {//所有用户的id存入到状态中,自动去重userIdState.put(element.getUserId(),"a");}Iterator<Long> iterator = userIdState.keys().iterator();ArrayList<Long> arrayList = new ArrayList<>();while(iterator.hasNext()){arrayList.add(iterator.next());}out.collect(context.window() + " " + arrayList.size());}}).print();env.execute();}}
![day07[Flink流处理高阶编程实战] - 图2](/uploads/projects/liuye-6lcqc@ddtw8t/70c195ab425faffa01d4f600ec0212c7.png)
第二章.电商数据分析
电商平台的用户行为频繁且比较复杂,系统上线运行一段时间后,可以收集到大量的用户行为数据,进而利用大数据技术进行深入挖掘和分析,得到感兴趣的商业指标并增强对风险的控制电商用户行为数据多样,整体可以分为用户行为数据(日志数据)和业务行为数据两大类用户的行为习惯数据包括了用户的登录方式,上线的时间点及时长,点击和浏览页面,页面停留时间以及页面跳转等等,我们可以从中进行流量统计和热门商品的统计,也可以深入挖掘用户的特征,这些数据往往可以从web服务器日志中直接读取到二业务行为数据就是用户在电商平台中针对每个业务(通常是某个具体商品)所做的具体操作,我们一般会在业务系统中相应的位置埋点,然后收集日志进行分析
1.数据准备
2.实时热门商品统计
需求分析: 每隔30分钟输出最近一小时内点击量最多的前N个商品
- 最近一小时:窗口长度
- 每隔5分钟:窗口滑动步长
- 时间:使用event-time(事件时间)
这里准备一个javabean用来保存商品信息
package com.atguigu.flink.day07.pojo;import lombok.AllArgsConstructor;import lombok.Data;import lombok.NoArgsConstructor;@Data@AllArgsConstructor@NoArgsConstructorpublic class HotItem {private Long itemId;private Long windowEndTime;private Long count;}
具体实现代码
package com.atguigu.flink.day07;import com.atguigu.flink.day07.pojo.HotItem;import com.atguigu.flink.day07.pojo.UserBehavior;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.AggregateFunction;import org.apache.flink.api.common.state.ListState;import org.apache.flink.api.common.state.ListStateDescriptor;import org.apache.flink.api.common.state.MapState;import org.apache.flink.api.common.state.MapStateDescriptor;import org.apache.flink.api.common.typeinfo.Types;import org.apache.flink.configuration.Configuration;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.KeyedProcessFunction;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows;import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import java.time.Duration;import java.util.ArrayList;import java.util.Iterator;/*** 需求:每隔30分钟输出最近一小时内点击量最多的前N个商品*/public class $03_TopN {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.readTextFile("input/UserBehavior.csv").map(line -> {String[] data = line.split(",");return new UserBehavior(Long.parseLong(data[0]),Long.parseLong(data[1]),Integer.parseInt(data[2]),data[3], Long.parseLong(data[4] ) * 1000);}).assignTimestampsAndWatermarks(WatermarkStrategy.<UserBehavior>forBoundedOutOfOrderness(Duration.ofSeconds(3)).withTimestampAssigner((ub,ts)-> ub.getTimestamp())).filter(ub -> "pv".equals(ub.getBehavior())).keyBy(ub -> ub.getItemId()).window(SlidingEventTimeWindows.of(Time.hours(1),Time.minutes(30))).aggregate(new AggregateFunction<UserBehavior, Long, Long>() {@Overridepublic Long createAccumulator() {return 0L;}@Overridepublic Long add(UserBehavior value, Long accumulator) {return accumulator + 1;}@Overridepublic Long getResult(Long accumulator) {return accumulator;}@Overridepublic Long merge(Long a, Long b) {return null;}},new ProcessWindowFunction<Long, HotItem, Long, TimeWindow>() {@Overridepublic void process(Long aLong, Context context, Iterable<Long> elements, Collector<HotItem> out) throws Exception {out.collect(new HotItem(aLong,context.window().getEnd(),elements.iterator().next()));}}).keyBy(hotItem -> hotItem.getWindowEndTime()).process(new KeyedProcessFunction<Long, HotItem, String>() {private ListState<HotItem> hotItemState;@Overridepublic void open(Configuration parameters) throws Exception {hotItemState = getRuntimeContext().getListState(new ListStateDescriptor<HotItem>("hotItemState",HotItem.class));}@Overridepublic void processElement(HotItem value, Context ctx, Collector<String> out) throws Exception {//每来一条数据,把数据存入到状态,等到定时器触发的时候,进行排序if(!hotItemState.get().iterator().hasNext()){//现在是来了第一条数据:注册定时器ctx.timerService().registerEventTimeTimer(value.getWindowEndTime() + 5000L);}hotItemState.add(value);}@Overridepublic void onTimer(long timestamp, OnTimerContext ctx, Collector<String> out) throws Exception {Iterator<HotItem> iterator = hotItemState.get().iterator();ArrayList<HotItem> hotItems = new ArrayList<>();while(iterator.hasNext()){hotItems.add(iterator.next());}hotItems.sort((o1,o2)->o2.getCount().compareTo(o1.getCount()));StringBuilder stringBuilder = new StringBuilder();stringBuilder.append("-------------------");for (int i = 0,count = Math.min(3, hotItems.size());i<count;i++){stringBuilder.append(hotItems.get(i)).append("\t");}out.collect(stringBuilder.toString());}}).print();env.execute();}}
![day07[Flink流处理高阶编程实战] - 图3](/uploads/projects/liuye-6lcqc@ddtw8t/e9ecd73d2939dabff2b148f8a20e4fb4.png)
第三章.页面广告分析
1.数据准备
2.黑名单过滤
我们进行的点击量计算,同一用户的重复点击是会叠加计算的在实际场景中,同一用户确实可能反复点开同一个广告,这也说明了用户对广告更大的兴趣,但是如果用户在一段时间频繁的点击广告,这显然不是一个正常行为,有刷点击量的嫌疑所以我们可以对一段时间内(比如一天内)的用户点击行为进行约束,如果对同一个广告点击超过一定限额(比如100次)应该把该用户加入黑名单并报警,此后其点击行为不应该再统计
两个功能
- 告警:使用侧输出流
- 已经进入黑名单的用户的广告点击记录不再进行统计
package com.atguigu.flink.day07;import com.atguigu.flink.day07.pojo.AdsClickLog;import com.atguigu.flink.day07.pojo.HotItem;import com.atguigu.flink.day07.pojo.UserBehavior;import org.apache.commons.lang.math.LongRange;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.AggregateFunction;import org.apache.flink.api.common.functions.ReduceFunction;import org.apache.flink.api.common.state.*;import org.apache.flink.api.common.typeinfo.Types;import org.apache.flink.configuration.Configuration;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.KeyedProcessFunction;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import org.apache.flink.util.OutputTag;import java.text.SimpleDateFormat;import java.time.Duration;import java.util.ArrayList;import java.util.Iterator;/*** 需求:黑名单过滤(时间间隔为一天,广告点击次数为100次)*/public class $04_BlackList {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);SingleOutputStreamOperator<String> normal = env.readTextFile("input/AdClickLog.csv").map(line -> {String[] data = line.split(",");return new AdsClickLog(Long.parseLong(data[0]),Long.parseLong(data[1]),data[2],data[3],Long.parseLong(data[4]));}).assignTimestampsAndWatermarks(WatermarkStrategy.<AdsClickLog>forBoundedOutOfOrderness(Duration.ofSeconds(3)).withTimestampAssigner((log, ts) -> log.getTimestamp())).keyBy(log -> log.getUserId() + ":" + log.getAdsId()).process(new KeyedProcessFunction<String, AdsClickLog, String>() {private ValueState<Boolean> warnState;private ReducingState<Long> countState;//定义一个状态存储年月日,如果新的数据和状态一致则是同一天,否则就是第二天private ValueState<String> yesterdayState;@Overridepublic void open(Configuration parameters) throws Exception {countState = getRuntimeContext().getReducingState(new ReducingStateDescriptor<Long>("countState",new ReduceFunction<Long>() {@Overridepublic Long reduce(Long value1, Long value2) throws Exception {return value1 + value2;}},Types.LONG));warnState = getRuntimeContext().getState(new ValueStateDescriptor<Boolean>("warnState", Types.BOOLEAN));yesterdayState = getRuntimeContext().getState(new ValueStateDescriptor<String>("yesterdayState", Types.STRING));}@Overridepublic void processElement(AdsClickLog value, Context ctx, Collector<String> out) throws Exception {String today = new SimpleDateFormat("yyyyMMdd").format(value.getTimestamp());String yesterday = yesterdayState.value();//如果到了第二天,清空状态点击量和黑名单重新计算if (!today.equals(yesterday)) {countState.clear();warnState.clear();yesterdayState.update(today);}if (warnState.value() == null) {//说明没有预警,需要对用户的操作进行统计countState.add(1L);}String msg = value.getUserId() + "对广告:" + value.getAdsId() + "的点击量:" + countState.get();if (countState.get() > 99) {if (warnState.value() == null) {ctx.output(new OutputTag<String>("black") {}, msg + ",超过了阈值100,加入黑名单");warnState.update(true);}} else {out.collect(msg);}}});normal.print("正常");normal.getSideOutput(new OutputTag<String>("black"){}).print("black");env.execute();}}
![day07[Flink流处理高阶编程实战] - 图4](/uploads/projects/liuye-6lcqc@ddtw8t/8493c4447e2c4b2093971a2a0f9932b0.png)
第四章.恶意登录监控
对于网站而言,用户登录并不是频繁的业务操作,如果一个用户短时间内频繁登录失败,就有可能是出现了程序的恶意攻击,比如密码暴力破解因此,我们考虑,应该对用户的登录失败动作进行统计,具体来说,如果一个同一用户在2s之内连续两次登录失败,就认为存在恶意登录的风险,输出相关的信息进行报警提示,这是电商网站,也是几乎所有网站风控的基本一环
1.数据准备
2.封装数据的JavaBean
package com.atguigu.flink.day07.pojo;import lombok.AllArgsConstructor;import lombok.Data;import lombok.NoArgsConstructor;@Data@NoArgsConstructor@AllArgsConstructorpublic class LoginEvent {private Long userId;private String ip;private String eventType;private Long eventTime;}
3.具体实现代码
package com.atguigu.flink.day07;import com.atguigu.flink.day07.pojo.LoginEvent;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.windowing.windows.GlobalWindow;import org.apache.flink.util.Collector;import java.util.ArrayList;import java.util.Iterator;/*** 需求:恶意登录监控(同一用户在2s之内连续两次登录失败)*/public class $05_Login {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);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);}).keyBy(log->log.getUserId()).countWindow(2,1).process(new ProcessWindowFunction<LoginEvent, String, Long, GlobalWindow>() {@Overridepublic void process(Long aLong, Context context, Iterable<LoginEvent> elements, Collector<String> out) throws Exception {Iterator<LoginEvent> iterator = elements.iterator();ArrayList<LoginEvent> loginEvents = new ArrayList<>();while(iterator.hasNext()){loginEvents.add(iterator.next());}if(loginEvents.size() ==2){//如果长度为2才需要进行计算LoginEvent event1 = loginEvents.get(0);LoginEvent event2 = loginEvents.get(1);String eventType1 = event1.getEventType();String eventType2 = event2.getEventType();Long time1 = event1.getEventTime();Long time2 = event2.getEventTime();if("fail".equals(eventType1)&& "fail".equals(eventType2)&& Math.abs(time1 -time2)<=2000){out.collect(aLong + "在恶意登录,请注意....");}}}}).print();env.execute();}}
![day07[Flink流处理高阶编程实战] - 图5](/uploads/projects/liuye-6lcqc@ddtw8t/a3ae7cf21360e282f8e00bfed7092dc8.png)
![day07[Flink流处理高阶编程实战] - 图6](/uploads/projects/liuye-6lcqc@ddtw8t/8bddebfeed4d8fa6ffd3c7a7a33c94e2.png)
第五章.订单支付实时监控
在电商网站中,订单的支付作为直接与营销收入挂钩的一环,在业务流程中非常重要对于订单而言,为了正确控制业务过程,也为了增加用户的支付意愿,网站一般会设置一个支付失效时间,超过一段时间不支付的订单就会被取消另外,对于订单的支付,我们还应保证用户支付的正确性,这可以通过第三方支付平台的交易数据来做一个实时对账
1.数据准备
2.封装数据的JavaBean
package com.atguigu.flink.day07.pojo;import lombok.AllArgsConstructor;import lombok.Data;import lombok.NoArgsConstructor;@Data@AllArgsConstructor@NoArgsConstructorpublic class OrderEvent {private Long orderId;private String eventType;private String txId;private Long eventTime;}
3.具体代码实现
package com.atguigu.flink.day07;import com.atguigu.flink.day07.pojo.OrderEvent;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.state.ValueState;import org.apache.flink.api.common.state.ValueStateDescriptor;import org.apache.flink.configuration.Configuration;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.windowing.assigners.EventTimeSessionWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import java.time.Duration;import java.util.ArrayList;import java.util.Iterator;/*** 需求:订单支付实时监控*/public class $06_Order {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);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()).window(EventTimeSessionWindows.withGap(Time.minutes(45))).process(new ProcessWindowFunction<OrderEvent, String, Long, TimeWindow>() {private ValueState<OrderEvent> creatState;@Overridepublic void open(Configuration parameters) throws Exception {creatState = getRuntimeContext().getState(new ValueStateDescriptor<OrderEvent>("creatState",OrderEvent.class));}@Overridepublic void process(Long key, Context context, Iterable<OrderEvent> elements, Collector<String> out) throws Exception {/*** 如果订单支付和创建都正常情况:则集合中有两个元素* 如果只有一个create,则表示没有key,或者pay的时间超过了45分钟*/Iterator<OrderEvent> iterator = elements.iterator();ArrayList<OrderEvent> orderEvents = new ArrayList<>();while (iterator.hasNext()) {orderEvents.add(iterator.next());}if (orderEvents.size() == 2) {out.collect("订单:" + key + "正常");} else {//如果长度不是2,则一定是1OrderEvent event = orderEvents.get(0);if ("create".equals(event.getEventType())) {//这个session窗口只有create,没有pay,要么真的没有,要么超时creatState.update(event);} else {if (creatState.value() != null) {//表示create来过但是pay晚于create超过了45分钟,他们没有在一个窗口中out.collect("订单:" + key + "超时支付,请检查系统漏洞");} else {out.collect("订单:" + key + " 只有pay没有create, 请检测系统漏洞");}}}}}).print();env.execute();}}
![day07[Flink流处理高阶编程实战] - 图7](/uploads/projects/liuye-6lcqc@ddtw8t/3d3f7273206c9d28af19895f142f8314.png)
