练习:Flink面试题


问题描述:

APP用户点击日志,列名分别为时间,用户ID,产品代号,点击的功能代号,邮箱,省市,耗时,参数详情。需使用flink批处理进行数据清洗及开窗统计,样例数据如下:

data

说明:

  1. 数据的列分隔符为逗号,详情参数为json
  2. 数据行中存在脏数据

环境:

  1. 机器可联网,笔试机器的桌面上有idea开发环境,flink相关依赖需自己引入(需使用flink 1.11.1 以上版本)
  2. 样例数据存放在

链接: https://pan.baidu.com/s/18n4PjyXHsrXwWz6rdBzKIw?pwd=ct6c 

提取码: ct6c

  1. 注意:笔试过程中避免重启机器,否则答题过程可能被还原
  2. 考试时间为2个小时

pom

<?xml version="1.0" encoding="UTF-8"?>

         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    4.0.0

    org.example
    FlinkTrue
    1.0-SNAPSHOT
    
        UTF-8
        1.8
        1.8
        1.13.0
        3.1.3
        2.12
        2.11
    

    

        
        
            org.apache.flink
            flink-core
            ${flink.version}
        
        
            org.apache.flink
            flink-java
            ${flink.version}
        
        
            org.apache.flink
            flink-streaming-java_${scala.version}
            ${flink.version}
        
        
            org.apache.flink
            flink-clients_${scala.version}
            ${flink.version}
        

        
            org.projectlombok
            lombok
            1.16.22
        


        
        
            redis.clients
            jedis
            3.3.0
        

        
        
            com.alibaba
            fastjson
            1.2.60
        


        
        
            org.apache.flink
            flink-table-api-java-bridge_${scala.version}
            ${flink.version}
        
        
            org.apache.flink
            flink-table-planner-blink_${scala.version}
            ${flink.version}
        
        
            org.apache.flink
            flink-table-common
            ${flink.version}
        

        
        
            org.apache.flink
            flink-cep-scala_${scala.version}
            ${flink.version}
        

        
        
            org.apache.flink
            flink-csv
            ${flink.version}
        

        
        
            org.apache.flink
            flink-connector-kafka_${scala.version}
            ${flink.version}
        

        
        
            org.apache.flink
            flink-connector-filesystem_${scala.binary.version}
            1.4.2
        
        
            org.apache.flink
            flink-hadoop-compatibility_${scala.version}
            ${flink.version}
        
        
            org.apache.hadoop
            hadoop-client
            ${hadoop.version}
        
        
            org.apache.hadoop
            hadoop-common
            ${hadoop.version}
        
        
            org.apache.hadoop
            hadoop-hdfs
            ${hadoop.version}
        

        
        
            org.apache.flink
            flink-jdbc_${scala.version}
            1.9.2
        
        
            mysql
            mysql-connector-java
            5.1.38
        

        

        
            org.apache.flink
            flink-hbase_${scala.version}
            1.8.1
        

        
            org.apache.hbase
            hbase-client
            2.4.3
        

        
        
            org.apache.flink
            flink-connector-jdbc_${scala.version}
            1.13.0
        
        
            ru.yandex.clickhouse
            clickhouse-jdbc
            0.2.4
        
        
        
            com.google.guava
            guava
            30.1.1-jre
        

        
        
            org.apache.flink
            flink-connector-jdbc_${scala.version}
            ${flink.version}
        
        
            ru.yandex.clickhouse
            clickhouse-jdbc
            0.2.4
            
                
                    com.fasterxml.jackson.core
                    jackson-databind
                
                
                    com.fasterxml.jackson.core
                    jackson-core
                
            
        

        
        
            org.apache.bahir
            flink-connector-redis_2.11
            1.0
        

        
        
            redis.clients
            jedis
            3.3.0
        

        
        
            org.apache.flink
            flink-connector-elasticsearch7_${scala.version}
            1.10.1
        

    

bean

import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;


@Data
@AllArgsConstructor
@NoArgsConstructor
public class Log {
    //时间
    private Long time   ;
    // 用户ID
    private String    uid ;
    // 产品代号
    private String    sid ;
    // 点击的功能代号
    private String    exeid ;
    // 邮箱
    private String    email ;
    // 省市
    private String    province ;
    // 耗时
    private String    spend ;
    // 参数详情
    private String    detail ;
}
import lombok.AllArgsConstructor;
import lombok.NoArgsConstructor;

@lombok.Data
@NoArgsConstructor
@AllArgsConstructor
public class Data {
    private Long time;
//    private Integer isman;
    private String level;
//    private Integer ts;
//    private Integer ver;
    private Integer count;
}

test

import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.AggregateFunction;
import org.apache.flink.api.common.functions.FilterFunction;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.functions.ReduceFunction;
import org.apache.flink.api.java.ExecutionEnvironment;
import org.apache.flink.api.java.operators.DataSource;
import org.apache.flink.streaming.api.TimeCharacteristic;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
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.timestamps.BoundedOutOfOrdernessTimestampExtractor;
import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
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.text.SimpleDateFormat;
import java.time.Duration;

public class FlinkTrue {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);

//        DataSource data = env.readTextFile(FlinkTrue.class.getClassLoader().getResource("simple.etl.csv").getPath());
        DataStreamSource data = env.readTextFile("C:\\Users\\liuyuan\\Desktop\\实训3\\FlinkTrue\\src\\main\\resources\\simple.etl.csv");
        SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
        env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
        SingleOutputStreamOperator maped = data.map(new MapFunction() {
            @Override
            public Log map(String s) throws Exception {
                String[] split = s.split(",");
                Log log = new Log();
                if (split.length >= 11) {
                    log = new Log(sdf.parse(split[0]).getTime(), split[1], split[2], split[3], split[4], split[5], split[6], split[7]+","+split[8]+","+split[9]+","+split[10]);
                }
                return log;
            }
        });
        SingleOutputStreamOperator map = maped.filter(new FilterFunction() {
            @Override
            public boolean filter(Log log) throws Exception {
                return !String.valueOf(log.getTime()).equals("null");
            }
        });

        SingleOutputStreamOperator mapData = map.map(new MapFunction() {
            @Override
            public Data map(Log log) throws Exception {
                String detail = log.getDetail();
                String s = detail.replaceAll("[\"]", "");
                JSONObject jsonObject = JSON.parseObject(s);
                return new Data(log.getTime(),String.valueOf(jsonObject.get("level")),0);
            }
        });

   /*     SingleOutputStreamOperator stream = mapData.assignTimestampsAndWatermarks(
                new BoundedOutOfOrdernessTimestampExtractor(Time.seconds(3)) {
            @Override
            public long extractTimestamp(Data data1) {
                                 return data1.getTime();
                             }
      });*/
        SingleOutputStreamOperator stream = mapData.assignTimestampsAndWatermarks(WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5))
                .withTimestampAssigner(new SerializableTimestampAssigner() {
                    @Override
                    public long extractTimestamp(Data data, long l) {
                        return data.getTime();
                    }
                }));


        SingleOutputStreamOperator aggregate = stream.keyBy(d -> d.getLevel())
                .window(TumblingEventTimeWindows.of(Time.seconds(30)))
                .aggregate(new AggregateFunction() {
                    @Override
                    public Integer createAccumulator() {
                        return 0;
                    }

                    @Override
                    public Integer add(Data data, Integer integer) {
                        data.setCount(integer + 1);
                        return data.getCount();
                    }

                    @Override
                    public Integer getResult(Integer integer) {
                        return integer;
                    }

                    @Override
                    public Integer merge(Integer integer, Integer acc1) {
                        return null;
                    }
                }, new ProcessWindowFunction() {

                    @Override
                    public void process(String s, Context context, Iterable iterable, Collector collector) throws Exception {
                        long end = context.window().getEnd();
                        long start = context.window().getStart();
                        Integer next = iterable.iterator().next();
                        collector.collect("结束时间:"+start+",结束时间:"+end+ ",level:"+s+",次数:"+next);
                    }
                });

        aggregate.print();


        //1.对数据进行必要的丢弃或修复
        //2.以样例数据中第1列时间作为事件时间
        //3.编写3s的窗口,统计详细参数里各level的出现次数
        //4.设定水位线为5s,过滤部分延迟数据

        //5.不允许使用过时方法·
        //6.请使用java语言实现
        //控制台输出:窗口时间+level+次数,样例如下 timewindow: key: count:

        env.execute();

    }
}

评价标准:

考察范围

数据丢弃及Json处理

结果正确,代码规范

结果正确,代码规范但存在瑕疵

不符合

Flink窗口使用

结果正确,代码规范

结果正确,代码规范但存在瑕疵

不符合

Flink水位线使用

结果正确,代码规范

结果正确,代码规范但存在瑕疵

不符合