第一章.Flink流处理核心编程

1.transform

一.简单滚动聚合算子

常见的滚动聚合算子 sum, min, max,minBy,maxBy

作用:KeyedStream的每一个支流做聚合,执行完成后,会将聚合的结果合成一个流返回,所以结果都是DataStream

参数:如果流中存储的是POJO或者scala的样例类, 参数使用字段名. 如果流中存储的是元组, 参数就是位置(基于0…)

返回:KeyedStream -> SingleOutputStreamOperator

  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.java.functions.KeySelector;
  4. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import java.util.ArrayList;
  7. /**
  8. * 简单滚动聚合算子
  9. */
  10. public class $01_RollingAggregationFunction {
  11. public static void main(String[] args) throws Exception {
  12. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  13. env.setParallelism(1);
  14. ArrayList<WaterSensor> waterSensors = new ArrayList<>();
  15. waterSensors.add(new WaterSensor("sensor_1", 1607527992000L, 20));
  16. waterSensors.add(new WaterSensor("sensor_1", 1607527994000L, 50));
  17. waterSensors.add(new WaterSensor("sensor_1", 1607527996000L, 50));
  18. waterSensors.add(new WaterSensor("sensor_2", 1607527993000L, 10));
  19. waterSensors.add(new WaterSensor("sensor_2", 1607527995000L, 30));
  20. DataStreamSource<WaterSensor> waterSensorDS = env.fromCollection(waterSensors);
  21. waterSensorDS
  22. .keyBy(new KeySelector<WaterSensor, String>() {
  23. @Override
  24. public String getKey(WaterSensor value) throws Exception {
  25. return value.getId();
  26. }
  27. })
  28. /**
  29. * maxBy和minBy可以指定当出现相同值的时候,其他字段是否取第一个.
  30. * true表示取第一个, false表示取与最大值(最小值)同一行的.
  31. * 5.3.10reduce
  32. */
  33. //.maxBy("vc",Boolean.FALSE).print();
  34. //.sum("vc").print();
  35. .min("vc").print();
  36. env.execute();
  37. }
  38. }

二.reduce

一个分组数据流的聚合操作,合并当前的元素和上次聚合的结果,产生一个新的值,返回的流中包含每一次聚合的结果,而不是只返回最后一次聚合的最终结果

为什么还要把中间值也保存下来,考虑流式数据的特点:没有终点,也就没有最终的概念了,任何一个中间的聚合结果都是值

注意:聚合后结果的类型,必须和原来流中元素的类型保持一致

  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.api.common.functions.ReduceFunction;
  5. import org.apache.flink.api.java.functions.KeySelector;
  6. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  7. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  8. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  9. import java.util.ArrayList;
  10. /**
  11. * reduce
  12. * 1.keyby之后调用
  13. * 2.每个分组的第一条数据来的时候,不会调用reduce方法,直接传递给下游
  14. * 3.value1是上一次的计算结果,value2是当前来的数据
  15. */
  16. public class $02_ReduceFunction {
  17. public static void main(String[] args) throws Exception {
  18. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  19. env.setParallelism(1);
  20. SingleOutputStreamOperator<WaterSensor> sensorDS = env
  21. .socketTextStream("hadoop102", 9999)
  22. .map(new MapFunction<String, WaterSensor>() {
  23. @Override
  24. public WaterSensor map(String value) throws Exception {
  25. String[] line = value.split(",");
  26. return new WaterSensor(
  27. line[0],
  28. Long.parseLong(line[1]),
  29. Integer.parseInt(line[2])
  30. );
  31. }
  32. });
  33. sensorDS
  34. .keyBy(sensor -> sensor.getId())
  35. .reduce(new ReduceFunction<WaterSensor>() {
  36. @Override
  37. public WaterSensor reduce(WaterSensor value1, WaterSensor value2) throws Exception {
  38. return new WaterSensor(
  39. value1.getId(),
  40. System.currentTimeMillis(),
  41. value1.getVc() + value2.getVc()
  42. );
  43. }
  44. })
  45. .print();
  46. env.execute();
  47. }
  48. }

三.process

process算子在flink算是一个比较底层的算子,很多类型的流上都可以调用,可以从流中获取更多的信息(不仅仅是数据本身)

  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
  7. import org.apache.flink.util.Collector;
  8. /**
  9. * process
  10. */
  11. public class $03_ProcessFunction {
  12. public static void main(String[] args) throws Exception {
  13. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  14. env.setParallelism(1);
  15. SingleOutputStreamOperator<WaterSensor> sensorDS = env
  16. .socketTextStream("hadoop102", 9999)
  17. .map(new MapFunction<String, WaterSensor>() {
  18. @Override
  19. public WaterSensor map(String value) throws Exception {
  20. String[] line = value.split(",");
  21. return new WaterSensor(
  22. line[0],
  23. Long.parseLong(line[1]),
  24. Integer.parseInt(line[2])
  25. );
  26. }
  27. });
  28. sensorDS
  29. .keyBy(sensor -> sensor.getId())
  30. .process(new KeyedProcessFunction<String, WaterSensor, String>() {
  31. @Override
  32. public void processElement(WaterSensor value, Context ctx, Collector<String> out) throws Exception {
  33. out.collect(ctx.getCurrentKey());
  34. }
  35. })
  36. .print();
  37. env.execute();
  38. }
  39. }

四.对流重新分区的几个算子

  • keyBy: 先按照key分组,按照key的双重hash来选择后面的分区
  • shuffle: 对流中的元素随机分区
  • rebalance : 对流中的元素平均分布到每个区,当处理倾斜数据的时候,进行性能优化
  • rescale: 同rebalance一样,也是平均循环的分布数据,但是要比rebalance更高效,因为rescale不需要通过网络,完全走的管道

2.sink

一.KafkaSink

  1. pom.xml
  1. <dependency>
  2. <groupId>org.apache.flink</groupId>
  3. <artifactId>flink-connector-kafka_2.12</artifactId>
  4. <version>1.13.1</version>
  5. </dependency>
  1. 代码
  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer;
  7. import org.apache.flink.streaming.connectors.kafka.KafkaSerializationSchema;
  8. import org.apache.kafka.clients.producer.ProducerRecord;
  9. import javax.annotation.Nullable;
  10. import java.nio.charset.StandardCharsets;
  11. import java.util.Properties;
  12. /**
  13. * 将数据写出到kafka
  14. * 生产者:
  15. * ack:
  16. * buffersize: 默认 32M
  17. * batchsize: 默认 16K
  18. * linger.ms: 默认 0ms
  19. *
  20. * 注意事项:
  21. * 1.可以指定一些生产者参数:ack.缓冲区,批次
  22. * 2.默认分区器 FlinkFixedPartitioner
  23. * 3.可以指定一致性级别: 至少一次,精准一次
  24. */
  25. public class $04_KafkaSink {
  26. public static void main(String[] args) throws Exception {
  27. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  28. env.setParallelism(1);
  29. SingleOutputStreamOperator<WaterSensor> sensorDS = env
  30. .socketTextStream("hadoop162", 9999)
  31. .map(new MapFunction<String, WaterSensor>() {
  32. @Override
  33. public WaterSensor map(String value) throws Exception {
  34. String[] line = value.split(",");
  35. return new WaterSensor(
  36. line[0],
  37. Long.parseLong(line[1]),
  38. Integer.parseInt(line[2])
  39. );
  40. }
  41. });
  42. Properties prop = new Properties();
  43. prop.put("bootstrap.servers","hadoop162:9092,hadoop163:9092,hadoop164:9092");
  44. sensorDS
  45. .addSink(new FlinkKafkaProducer<WaterSensor>(
  46. "default",
  47. new KafkaSerializationSchema<WaterSensor>() {
  48. @Override
  49. public ProducerRecord<byte[], byte[]> serialize(WaterSensor waterSensor, @Nullable Long aLong) {
  50. return new ProducerRecord<byte[], byte[]>("waterSensor",waterSensor.toString().getBytes(StandardCharsets.UTF_8));
  51. }
  52. },
  53. prop,
  54. FlinkKafkaProducer.Semantic.EXACTLY_ONCE
  55. ));
  56. env.execute();
  57. }
  58. }

二.Redis Sink

  1. pom.xml
  1. <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-connector-redis -->
  2. <dependency>
  3. <groupId>org.apache.flink</groupId>
  4. <artifactId>flink-connector-redis_2.11</artifactId>
  5. <version>1.1.5</version>
  6. </dependency>
  1. 代码
  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer;
  7. import org.apache.flink.streaming.connectors.kafka.KafkaSerializationSchema;
  8. import org.apache.flink.streaming.connectors.redis.RedisSink;
  9. import org.apache.flink.streaming.connectors.redis.common.config.FlinkJedisPoolConfig;
  10. import org.apache.flink.streaming.connectors.redis.common.mapper.RedisCommand;
  11. import org.apache.flink.streaming.connectors.redis.common.mapper.RedisCommandDescription;
  12. import org.apache.flink.streaming.connectors.redis.common.mapper.RedisMapper;
  13. import org.apache.kafka.clients.producer.ProducerRecord;
  14. import javax.annotation.Nullable;
  15. import java.nio.charset.StandardCharsets;
  16. import java.util.Properties;
  17. /**
  18. * 将数据写入到redis
  19. */
  20. public class $05_RedisSink {
  21. public static void main(String[] args) throws Exception {
  22. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  23. env.setParallelism(1);
  24. SingleOutputStreamOperator<WaterSensor> sensorDS = env
  25. .socketTextStream("hadoop162", 9999)
  26. .map(new MapFunction<String, WaterSensor>() {
  27. @Override
  28. public WaterSensor map(String value) throws Exception {
  29. String[] line = value.split(",");
  30. return new WaterSensor(
  31. line[0],
  32. Long.parseLong(line[1]),
  33. Integer.parseInt(line[2])
  34. );
  35. }
  36. });
  37. FlinkJedisPoolConfig jedisPoolConfig = new FlinkJedisPoolConfig.Builder()
  38. .setHost("hadoop162")
  39. .setPort(6379)
  40. .build();
  41. RedisSink<WaterSensor> redisSinkFunction = new RedisSink<>(
  42. jedisPoolConfig,
  43. new RedisMapper<WaterSensor>() {
  44. @Override
  45. public RedisCommandDescription getCommandDescription() {
  46. //指定命令:往Hash结构插入数据(也可以是其他命令)
  47. return new RedisCommandDescription(RedisCommand.HSET, "waterSensor");
  48. }
  49. /**
  50. * 从数据读取hash结构的key
  51. * @param waterSensor
  52. * @return
  53. */
  54. @Override
  55. public String getKeyFromData(WaterSensor waterSensor) {
  56. return waterSensor.getId();
  57. }
  58. /**
  59. * 从数据提取,hash结构的value
  60. * @param waterSensor
  61. * @return
  62. */
  63. @Override
  64. public String getValueFromData(WaterSensor waterSensor) {
  65. return waterSensor.getVc().toString();
  66. }
  67. }
  68. );
  69. sensorDS.addSink(redisSinkFunction);
  70. env.execute();
  71. }
  72. }

三.ElasticsearchSink

  1. pom.xml
  1. <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-connector-elasticsearch6 -->
  2. <dependency>
  3. <groupId>org.apache.flink</groupId>
  4. <artifactId>flink-connector-elasticsearch6_2.12</artifactId>
  5. <version>1.13.1</version>
  6. </dependency>
  1. 代码
  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.api.common.functions.RuntimeContext;
  5. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  6. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  7. import org.apache.flink.streaming.connectors.elasticsearch.ElasticsearchSinkFunction;
  8. import org.apache.flink.streaming.connectors.elasticsearch.RequestIndexer;
  9. import org.apache.flink.streaming.connectors.elasticsearch6.ElasticsearchSink;
  10. import org.apache.http.HttpHost;
  11. import org.elasticsearch.action.index.IndexRequest;
  12. import org.elasticsearch.client.Requests;
  13. import java.util.ArrayList;
  14. import java.util.HashMap;
  15. /**
  16. * 将数据写入到es
  17. * es: http://hadoop162:9200/watersensor/_search
  18. * kibana: GET /watersensor/_search
  19. */
  20. public class $06_ElasticsearchSink {
  21. public static void main(String[] args) throws Exception {
  22. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  23. env.setParallelism(1);
  24. SingleOutputStreamOperator<WaterSensor> sensorDS = env
  25. .socketTextStream("hadoop162", 9999)
  26. .map(new MapFunction<String, WaterSensor>() {
  27. @Override
  28. public WaterSensor map(String value) throws Exception {
  29. String[] line = value.split(",");
  30. return new WaterSensor(
  31. line[0],
  32. Long.parseLong(line[1]),
  33. Integer.parseInt(line[2])
  34. );
  35. }
  36. });
  37. ArrayList<HttpHost> httpHosts = new ArrayList<>();
  38. httpHosts.add(new HttpHost("hadoop162",9200));
  39. httpHosts.add(new HttpHost("hadoop163",9200));
  40. httpHosts.add(new HttpHost("hadoop164",9200));
  41. ElasticsearchSink.Builder<WaterSensor> esBuilder = new ElasticsearchSink.Builder<>(
  42. httpHosts,
  43. new ElasticsearchSinkFunction<WaterSensor>() {
  44. @Override
  45. public void process(WaterSensor waterSensor, RuntimeContext runtimeContext, RequestIndexer requestIndexer) {
  46. HashMap<String, String> dataMap = new HashMap<>();
  47. dataMap.put("data", waterSensor.toString());
  48. IndexRequest indexRequest = Requests.indexRequest()
  49. .index("watersensor")
  50. .type("_doc")
  51. .source(dataMap);
  52. requestIndexer.add(indexRequest);
  53. }
  54. }
  55. );
  56. //TODO 注意:这个仅仅是为了演示效果才设为1,生产环境不要改成1
  57. esBuilder.setBulkFlushMaxActions(1);
  58. ElasticsearchSink<WaterSensor> esSinkFunction = esBuilder.build();
  59. sensorDS.addSink(esSinkFunction);
  60. env.execute();
  61. }
  62. }

day03[Flink流处理核心编程实战] - 图1

四.自定义sink

需求: 使用flink读取socket数据到mysql

  1. 导入mysql依赖
  1. <dependency>
  2. <groupId>mysql</groupId>
  3. <artifactId>mysql-connector-java</artifactId>
  4. <version>5.1.49</version>
  5. </dependency>
  1. 建库建表
  1. create database test;
  2. use test;
  3. CREATE TABLE `sensor` (
  4. `id` varchar(20) NOT NULL,
  5. `ts` bigint(20) NOT NULL,
  6. `vc` int(11) NOT NULL,
  7. PRIMARY KEY (`id`,`ts`)
  8. ) ENGINE=InnoDB DEFAULT CHARSET=utf8;
  1. 代码
  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.api.common.functions.RuntimeContext;
  5. import org.apache.flink.configuration.Configuration;
  6. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  7. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  8. import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
  9. import org.apache.flink.streaming.connectors.elasticsearch.ElasticsearchSinkFunction;
  10. import org.apache.flink.streaming.connectors.elasticsearch.RequestIndexer;
  11. import org.apache.flink.streaming.connectors.elasticsearch6.ElasticsearchSink;
  12. import org.apache.http.HttpHost;
  13. import org.elasticsearch.action.index.IndexRequest;
  14. import org.elasticsearch.client.Requests;
  15. import java.sql.Connection;
  16. import java.sql.DriverManager;
  17. import java.sql.PreparedStatement;
  18. import java.util.ArrayList;
  19. import java.util.HashMap;
  20. /**
  21. * 自定义sink :将数据写出到kafka
  22. * 生产环境用官方的JDBC写法
  23. * https://nightlies.apache.org/flink/flink-docs-release-1.13/docs/connectors/datastream/jdbc/
  24. */
  25. public class $07_MySqlSink {
  26. public static void main(String[] args) throws Exception {
  27. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  28. env.setParallelism(1);
  29. SingleOutputStreamOperator<WaterSensor> sensorDS = env
  30. .socketTextStream("hadoop162", 9999)
  31. .map(new MapFunction<String, WaterSensor>() {
  32. @Override
  33. public WaterSensor map(String value) throws Exception {
  34. String[] line = value.split(",");
  35. return new WaterSensor(
  36. line[0],
  37. Long.parseLong(line[1]),
  38. Integer.parseInt(line[2])
  39. );
  40. }
  41. });
  42. sensorDS.addSink(new MySink());
  43. env.execute();
  44. }
  45. public static class MySink extends RichSinkFunction<WaterSensor>{
  46. private Connection connection;
  47. private PreparedStatement preparedStatement;
  48. @Override
  49. public void open(Configuration parameters) throws Exception {
  50. connection = DriverManager.getConnection("jdbc:mysql://hadoop162:3306/test", "root", "aaaaaa");
  51. preparedStatement = connection.prepareStatement("insert into sensor values (?,?,?)");
  52. }
  53. @Override
  54. public void close() throws Exception {
  55. if(preparedStatement != null){
  56. preparedStatement.close();
  57. }
  58. if(connection != null){
  59. connection.close();
  60. }
  61. }
  62. @Override
  63. public void invoke(WaterSensor value, Context context) throws Exception {
  64. preparedStatement.setString(1, value.getId());
  65. preparedStatement.setLong(2,value.getTs());
  66. preparedStatement.setInt(3,value.getVc());
  67. preparedStatement.execute();
  68. }
  69. }
  70. }

3.配置BATH执行模式

执行模式有3个选择可配

  1. STREAMING(默认)
  2. BATCH
  3. AUTOMATIC
  1. 有界数据用哪个STREAMING 和BATCH的区别
  2. STREAMING模式下, 数据是来一条输出一次结果.
  3. BATCH模式下, 数据处理完之后, 一次性输出结果.

命令行配置

bin/flink run -Dexecution.runtime-mode=BATCH …

第二章.Flink流处理核心编程实战

1.基于埋点日志数据的网络流量统计

一.网站总浏览量(PV)的统计

  1. 衡量网站流量一个最简单的指标,就是网站的页面浏览量(Page View,PV)。用户每次打开一个页面便记录1次PV,多次打开同一页面则浏览量累计。
  2. 一般来说,PV与来访者的数量成正比,但是PV并不直接决定页面的真实来访者数量,如同一个来访者通过不断的刷新页面,也可以制造出非常高的PV。接下来我们就用咱们之前学习的Flink算子来实现PV的统计
  1. 用于封装数据的javabean
  1. package com.atguigu.flink.day03.pojo;
  2. import lombok.AllArgsConstructor;
  3. import lombok.Data;
  4. import lombok.NoArgsConstructor;
  5. @Data
  6. @NoArgsConstructor
  7. @AllArgsConstructor
  8. public class UserBehavior {
  9. /**
  10. * 用户ID
  11. */
  12. private Long userId;
  13. /**
  14. * 商品ID
  15. */
  16. private Long itemId;
  17. /**
  18. * 品类ID
  19. */
  20. private Integer categoryId;
  21. /**
  22. * 用户的行为:pv、buy、cart、fav
  23. */
  24. private String behavior;
  25. /**
  26. * 时间戳
  27. */
  28. private Long timestamp;
  29. }
  1. 代码
  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day03.pojo.UserBehavior;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.streaming.api.functions.ProcessFunction;
  7. import org.apache.flink.util.Collector;
  8. /**
  9. * pv
  10. */
  11. public class $08_PV {
  12. public static void main(String[] args) throws Exception {
  13. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  14. env.setParallelism(2);
  15. //1.读取数据
  16. SingleOutputStreamOperator<UserBehavior> userBehavioirDS = env
  17. .readTextFile("input/UserBehavior.csv")
  18. .map(new MapFunction<String, UserBehavior>() {
  19. @Override
  20. public UserBehavior map(String value) throws Exception {
  21. String[] datas = value.split(",");
  22. return new UserBehavior(
  23. Long.parseLong(datas[0]),
  24. Long.parseLong(datas[1]),
  25. Integer.parseInt(datas[2]),
  26. datas[3],
  27. Long.parseLong(datas[4])
  28. );
  29. }
  30. });
  31. //2.处理数据
  32. userBehavioirDS
  33. .filter(user -> "pv".equals(user.getBehavior()))
  34. .process(new ProcessFunction<UserBehavior, Integer>() {
  35. int pvCount = 0;
  36. @Override
  37. public void processElement(UserBehavior value, Context ctx, Collector<Integer> out) throws Exception {
  38. pvCount++;
  39. out.collect(pvCount);
  40. }
  41. }).setParallelism(1)
  42. .print();
  43. env.execute();
  44. }
  45. }

二.网站独立访客数(UV)的统计

上一个案例中,我们统计的是所有用户对页面的所有浏览行为,也就是说,同一用户的浏览行为会被重复统计。而在实际应用中,我们往往还会关注,到底有多少不同的用户访问了网站,所以另外一个统计流量的重要指标是网站的独立访客数(Unique Visitor,UV)

对于UserBehavior数据源来说,我们直接可以根据userId来区分不同的用户.

代码

  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day03.pojo.UserBehavior;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.streaming.api.functions.ProcessFunction;
  7. import org.apache.flink.util.Collector;
  8. import java.util.HashSet;
  9. import java.util.Set;
  10. /**
  11. * uv
  12. */
  13. public class $09_UV {
  14. public static void main(String[] args) throws Exception {
  15. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  16. env.setParallelism(2);
  17. //1.读取数据
  18. SingleOutputStreamOperator<UserBehavior> userBehavioirDS = env
  19. .readTextFile("input/UserBehavior.csv")
  20. .map(new MapFunction<String, UserBehavior>() {
  21. @Override
  22. public UserBehavior map(String value) throws Exception {
  23. String[] datas = value.split(",");
  24. return new UserBehavior(
  25. Long.parseLong(datas[0]),
  26. Long.parseLong(datas[1]),
  27. Integer.parseInt(datas[2]),
  28. datas[3],
  29. Long.parseLong(datas[4])
  30. );
  31. }
  32. });
  33. //2.处理数据,提取出 用户ID,存在SET结构里,求SET的长度就是UV值
  34. userBehavioirDS
  35. .filter(user -> "pv".equals(user.getBehavior()))
  36. .process(new ProcessFunction<UserBehavior, Integer>() {
  37. Set<Long> uvSet = new HashSet<>();
  38. @Override
  39. public void processElement(UserBehavior value, Context ctx, Collector<Integer> out) throws Exception {
  40. //每来一条数据,提取出用户ID,存到Set里
  41. uvSet.add(value.getItemId());
  42. //将Set的元素个数(uv值),传递给下游
  43. out.collect(uvSet.size());
  44. }
  45. }).setParallelism(1)
  46. .print();
  47. env.execute();
  48. }
  49. }

2.市场营销商业指标统计分析

随着智能手机的普及,在如今的电商网站中已经有越来越多的用户来自移动端,相比起传统浏览器的登录方式,手机APP成为了更多用户访问电商网站的首选。对于电商企业来说,一般会通过各种不同的渠道对自己的APP进行市场推广,而这些渠道的统计数据(比如,不同网站上广告链接的点击量、APP下载量)就成了市场营销的重要商业指标。

一.APP市场推广统计 -分渠道

  1. 封装数据的javabean
  1. package com.atguigu.flink.day03.pojo;
  2. import lombok.AllArgsConstructor;
  3. import lombok.Data;
  4. import lombok.NoArgsConstructor;
  5. @Data
  6. @AllArgsConstructor
  7. @NoArgsConstructor
  8. public class MarketingUserBehavior {
  9. private Long userId;
  10. private String behavior;
  11. private String channel;
  12. private Long timestamp;
  13. }
  1. 具体实现代码
  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day03.pojo.MarketingUserBehavior;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.api.java.tuple.Tuple2;
  5. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  6. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  7. import org.apache.flink.streaming.api.functions.source.SourceFunction;
  8. import java.util.Arrays;
  9. import java.util.List;
  10. import java.util.Random;
  11. /**
  12. * APP市场推广统计-分渠道
  13. */
  14. public class $10_AppAnalysisByChannel {
  15. public static void main(String[] args) throws Exception {
  16. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  17. env.setParallelism(2);
  18. //1.读取数据
  19. DataStreamSource<MarketingUserBehavior> inputDS = env.addSource(new AppSource());
  20. //2.处理数据:按照不同渠道对不同行为进行统计-->对两个维度进行keyby
  21. inputDS
  22. .map(new MapFunction<MarketingUserBehavior, Tuple2<String,Integer>>() {
  23. @Override
  24. public Tuple2<String, Integer> map(MarketingUserBehavior value) throws Exception {
  25. return Tuple2.of(value.getChannel() +"_"+value.getBehavior(),1);
  26. }
  27. })
  28. .keyBy(t->t.f0)
  29. .sum(1)
  30. .print();
  31. env.execute();
  32. }
  33. public static class AppSource implements SourceFunction<MarketingUserBehavior>{
  34. static boolean flag = true;
  35. List<String> behaviorList = Arrays.asList("DOWNLOAD", "INSTALL", "UPDATE", "UNINSTALL");
  36. List<String> channelList = Arrays.asList("XIAOMI", "HUAWEI", "OPPO", "VIVO");
  37. @Override
  38. public void run(SourceContext<MarketingUserBehavior> ctx) throws Exception {
  39. Random random = new Random();
  40. while (flag) {
  41. ctx.collect(
  42. new MarketingUserBehavior(
  43. random.nextLong(),
  44. behaviorList.get(random.nextInt(behaviorList.size())),
  45. channelList.get(random.nextInt(channelList.size())),
  46. System.currentTimeMillis()
  47. )
  48. );
  49. Thread.sleep(1000);
  50. }
  51. }
  52. @Override
  53. public void cancel() {
  54. flag = false;
  55. }
  56. }
  57. }

二.APP市场推广统计 -不分渠道

  1. 实现代码
  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day03.pojo.MarketingUserBehavior;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.api.java.tuple.Tuple2;
  5. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  6. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  7. import org.apache.flink.streaming.api.functions.source.SourceFunction;
  8. import java.util.Arrays;
  9. import java.util.List;
  10. import java.util.Random;
  11. /**
  12. * APP市场推广统计-不分渠道
  13. */
  14. public class $11_AppAnalysisWithOutChannel {
  15. public static void main(String[] args) throws Exception {
  16. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  17. env.setParallelism(2);
  18. //1.读取数据
  19. DataStreamSource<MarketingUserBehavior> inputDS = env.addSource(new AppSource());
  20. //2.处理数据:按照不同渠道对不同行为进行统计-->对两个维度进行keyby
  21. inputDS
  22. .map(new MapFunction<MarketingUserBehavior, Tuple2<String,Integer>>() {
  23. @Override
  24. public Tuple2<String, Integer> map(MarketingUserBehavior value) throws Exception {
  25. return Tuple2.of(value.getBehavior(),1);
  26. }
  27. })
  28. .keyBy(t->t.f0)
  29. .sum(1)
  30. .print();
  31. env.execute();
  32. }
  33. public static class AppSource implements SourceFunction<MarketingUserBehavior>{
  34. static boolean flag = true;
  35. List<String> behaviorList = Arrays.asList("DOWNLOAD", "INSTALL", "UPDATE", "UNINSTALL");
  36. List<String> channelList = Arrays.asList("XIAOMI", "HUAWEI", "OPPO", "VIVO");
  37. @Override
  38. public void run(SourceContext<MarketingUserBehavior> ctx) throws Exception {
  39. Random random = new Random();
  40. while (flag) {
  41. ctx.collect(
  42. new MarketingUserBehavior(
  43. random.nextLong(),
  44. behaviorList.get(random.nextInt(behaviorList.size())),
  45. channelList.get(random.nextInt(channelList.size())),
  46. System.currentTimeMillis()
  47. )
  48. );
  49. Thread.sleep(1000);
  50. }
  51. }
  52. @Override
  53. public void cancel() {
  54. flag = false;
  55. }
  56. }
  57. }

3.各省市页面广告点击量实时统计

  1. 电商网站的市场营销商业指标中,除了自身的APP推广,还会考虑到页面上的广告投放(包括自己经营的产品和其它网站的广告)。所以广告相关的统计分析,也是市场营销的重要指标。<br /> 对于广告的统计,最简单也最重要的就是页面广告的点击量,网站往往需要根据广告点击量来制定定价策略和调整推广方式,而且也可以借此收集用户的偏好信息。更加具体的应用是,我们可以根据用户的地理位置进行划分,从而总结出不同省份用户对不同广告的偏好,这样更有助于广告的精准投放。
  1. 用于封装数据的javabean
  1. package com.atguigu.flink.day03.pojo;
  2. import lombok.AllArgsConstructor;
  3. import lombok.Data;
  4. import lombok.NoArgsConstructor;
  5. @Data
  6. @AllArgsConstructor
  7. @NoArgsConstructor
  8. public class AdsClickLog {
  9. private Long userId;
  10. private Long adId;
  11. private String province;
  12. private String city;
  13. private Long timestamp;
  14. }
  1. 实现代码
  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day03.pojo.AdsClickLog;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.api.java.functions.KeySelector;
  5. import org.apache.flink.api.java.tuple.Tuple2;
  6. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  7. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  8. import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
  9. import org.apache.flink.util.Collector;
  10. import java.util.*;
  11. /**
  12. * 各省市页面广告点击量实时统计
  13. *
  14. * 扩展知识:
  15. * 函数类中,定义的变量个数和作用范围
  16. * 1.函数类创建了几次?
  17. * 每个并行实例new一次
  18. * 2.类中定义的变量一共有几个?
  19. * 每个并行实例一份
  20. * 3.函数类的方法调用几次?
  21. * 大部分是,来一条数据调用一次方法
  22. * 也有一些是,来一批数据调用一次(比如窗口的全量函数)
  23. */
  24. public class $12_AdClick {
  25. public static void main(String[] args) throws Exception {
  26. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  27. env.setParallelism(2);
  28. //1.读取数据,转换格式
  29. SingleOutputStreamOperator<AdsClickLog> adClickDS = env
  30. .readTextFile("input/AdClickLog.csv")
  31. .map(new MapFunction<String, AdsClickLog>() {
  32. @Override
  33. public AdsClickLog map(String value) throws Exception {
  34. String[] datas = value.split(",");
  35. return new AdsClickLog(
  36. Long.parseLong(datas[0]),
  37. Long.parseLong(datas[1]),
  38. datas[2],
  39. datas[3],
  40. Long.parseLong(datas[4])
  41. );
  42. }
  43. });
  44. //2.处理数据
  45. adClickDS
  46. .keyBy(new KeySelector<AdsClickLog, Tuple2<String,Long>>() {
  47. @Override
  48. public Tuple2<String, Long> getKey(AdsClickLog value) throws Exception {
  49. return Tuple2.of(value.getProvince(),value.getAdId());
  50. }
  51. })
  52. .process(new KeyedProcessFunction<Tuple2<String, Long>, AdsClickLog, String>() {
  53. Map<Tuple2<String, Long>, Integer> map = new HashMap<>();
  54. @Override
  55. public void processElement(AdsClickLog value, Context ctx, Collector<String> out) throws Exception {
  56. Integer count = map.get(ctx.getCurrentKey());
  57. if(count == null){
  58. //说明当前的key不存在,说明这是第一条数据
  59. count = 1;
  60. }else{
  61. //说明当前的key存在,说明这不是第一条数据,正常+1就行
  62. count++;
  63. }
  64. map.put(ctx.getCurrentKey(),count);
  65. out.collect(map.toString());
  66. }
  67. })
  68. .print();
  69. env.execute();
  70. }
  71. }

4.订单支付实时监控

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

需求:来自两条流的订单交易匹配

  1. 对于订单支付事件,用户支付完成其实并不算完,我们还得确认平台账户上是否到账了。而往往这会来自不同的日志信息,所以我们要同时读入两条流的数据来做合并处理。
  1. 用于封装数据的javabean
  1. package com.atguigu.flink.day03.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. }
  1. package com.atguigu.flink.day03.pojo;
  2. import lombok.AllArgsConstructor;
  3. import lombok.Data;
  4. import lombok.NoArgsConstructor;
  5. @Data
  6. @AllArgsConstructor
  7. @NoArgsConstructor
  8. public class TxEvent {
  9. private String txId;
  10. private String payChannel;
  11. private Long eventTime;
  12. }
  1. 实现代码
  1. package com.atguigu.flink.day03;
  2. import com.atguigu.flink.day03.pojo.OrderEvent;
  3. import com.atguigu.flink.day03.pojo.TxEvent;
  4. import org.apache.flink.api.common.functions.MapFunction;
  5. import org.apache.flink.streaming.api.datastream.ConnectedStreams;
  6. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  7. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  8. import org.apache.flink.streaming.api.functions.co.CoProcessFunction;
  9. import org.apache.flink.util.Collector;
  10. import java.util.HashMap;
  11. import java.util.Map;
  12. /**
  13. *connect连接两条流的时候,如果需要做一些比对操作,要考虑数据的分布,
  14. * 需要比对的数据要在同一个分区里 ---->keyby(关联条件)
  15. */
  16. public class $13_OrderTxDetect {
  17. public static void main(String[] args) throws Exception {
  18. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  19. env.setParallelism(1);
  20. //1.读取数据,转换格式
  21. SingleOutputStreamOperator<OrderEvent> orderDS = env.readTextFile("input/OrderLog.csv")
  22. .map(new MapFunction<String, OrderEvent>() {
  23. @Override
  24. public OrderEvent map(String value) throws Exception {
  25. String[] datas = value.split(",");
  26. return new OrderEvent(
  27. Long.parseLong(datas[0]),
  28. datas[1],
  29. datas[2],
  30. Long.parseLong(datas[3])
  31. );
  32. }
  33. });
  34. SingleOutputStreamOperator<TxEvent> txDS = env.readTextFile("input/ReceiptLog.csv")
  35. .map(new MapFunction<String, TxEvent>() {
  36. @Override
  37. public TxEvent map(String value) throws Exception {
  38. String[] datas = value.split(",");
  39. return new TxEvent(
  40. datas[0],
  41. datas[1],
  42. Long.parseLong(datas[2])
  43. );
  44. }
  45. });
  46. //2.数据处理,需要做一个合流操作
  47. ConnectedStreams<OrderEvent, TxEvent> orderTxCS = orderDS.connect(txDS);
  48. ConnectedStreams<OrderEvent, TxEvent> orderTxKCS = orderTxCS.keyBy(order -> order.getTxId(), tx -> tx.getTxId());
  49. SingleOutputStreamOperator<String> resultDS = orderTxKCS.process(new CoProcessFunction<OrderEvent, TxEvent, String>() {
  50. Map<String, OrderEvent> orderMap = new HashMap<>();
  51. Map<String, TxEvent> txMap = new HashMap<>();
  52. /**
  53. * 处理的是业务系统的数据
  54. * @param value
  55. * @param ctx
  56. * @param out
  57. * @throws Exception
  58. */
  59. @Override
  60. public void processElement1(OrderEvent value, Context ctx, Collector<String> out) throws Exception {
  61. //进入这个方法,说明当前来的是业务系统的数据,先查一下,对方交易(数据)来了没有
  62. TxEvent txEvent = txMap.get(value.getTxId());
  63. if (txEvent == null) {
  64. //说明交易数据没来过,等他先把自己存起来
  65. orderMap.put(value.getTxId(), value);
  66. } else {
  67. //说明,对方来过==>匹配上,删除map中的数据
  68. out.collect("订单" + value.getOrderId() + "对账成功");
  69. txMap.remove(value.getTxId());
  70. }
  71. }
  72. /**
  73. * 处理的是交易系统的数据
  74. * @param value
  75. * @param ctx
  76. * @param out
  77. * @throws Exception
  78. */
  79. @Override
  80. public void processElement2(TxEvent value, Context ctx, Collector<String> out) throws Exception {
  81. //进入这个方法,说明当前来的是交易系统的数据,先查一下,对方业务(数据)来了没有
  82. OrderEvent orderEvent = orderMap.get(value.getTxId());
  83. if (orderEvent == null) {
  84. //说明业务数据没来过,等他先把自己存起来
  85. txMap.put(value.getTxId(), value);
  86. } else {
  87. //说明,对方来过==>匹配上,删除map中的数据
  88. out.collect("订单" + orderEvent.getOrderId() + "对账成功");
  89. orderMap.remove(value.getTxId());
  90. }
  91. }
  92. });
  93. resultDS.print();
  94. env.execute();
  95. }
  96. }