Flink Environment & Source


Environment

 Flink Job在提交执行计算时,需要首先建立和Flink框架之间的联系,也就指的是当前的flink运行环境,只有获取了环境信息,才能将task调度到不同的taskManager执行。而这个环境对象的获取方式相对比较简单

// 批处理环境
ExecutionEnvironment benv = ExecutionEnvironment.getExecutionEnvironment();
// 流式数据处理环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

Source

 Flink框架可以从不同的来源获取数据,将数据提交给框架进行处理, 我们将获取数据的来源称之为数据源(Source)。

准备工作

引入依赖


<dependency>
    <groupId>org.projectlombokgroupId>
    <artifactId>lombokartifactId>
    <version>1.18.16version>
    <scope>providedscope>
dependency>

创建实体类

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

/**
 * 水位 监控器 用于接收水位数据
 * id 传感器编号
 * ts 时间戳
 * vc 水位
 */
@Data
@NoArgsConstructor
@AllArgsConstructor
public class WaterSensor {

    public String id;
    public long ts;
    public Integer vc;
}

从集合中读取数据

一般情况下,可以将数据临时存储到内存中,形成特殊的数据结构后,作为数据源使用。这里的数据结构采用集合类型是比较普遍的。

import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.wdh01.bean.WaterSensor;

import java.util.Arrays;
import java.util.List;

/**
 * 从集合读取数据
 */
public class Flink01_Source_Collection {
    public static void main(String[] args) throws Exception {
        //1、获取执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(2);//若 并行度1 按顺序读取数据,否则不按顺序
        //2、准备集合
        List waterSensors = Arrays.asList(
                new WaterSensor("ws_001", 1577844001L, 45),
                new WaterSensor("ws_002", 1577844015L, 43),
                new WaterSensor("ws_003", 1577844020L, 42));
        //3、从集合读取数据
        DataStreamSource waterSensorDataStreamSource = env.fromCollection(waterSensors);
        //4、打印
        waterSensorDataStreamSource.print();
        //5、执行
        env.execute();
    }
}

从文件中读取数据

import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.wdh01.bean.WaterSensor;

/**
 * 从文件读取数据
 */
public class Flink02_Source_File {
    public static void main(String[] args) throws Exception {
        //1、获取执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);//若 并行度1 按顺序读取数据,否则不按顺序
        // 2、从文件读取数据
        DataStreamSource stringDataStreamSource = env.readTextFile("input/sensor");
        //3、转换 javabean 打印数据 new MapFunction
        stringDataStreamSource.map(new MapFunction() {
            public WaterSensor map(String value) throws Exception {
                String[] split = value.split(",");
                return new WaterSensor(split[0], Long.parseLong(split[1]), Integer.parseInt(split[2]));
            }
        }).print();
        //4、执行
        env.execute();
    }
}

说明

也可以从 hdfs 文件系统读取数据,
  1. 参数可以是目录也可以是文件
  2. 路径可以是相对路径也可以是绝对路径
  3. 相对路径是从系统属性user.dir获取路径: idea下是project的根目录, standalone模式下是集群节点根目录
  4. 也可以从hdfs目录下读取, 使用路径:hdfs://...., 由于Flink没有提供hadoop相关依赖, 需要pom中添加相关依赖:
<dependency>
    <groupId>org.apache.hadoopgroupId>
    <artifactId>hadoop-clientartifactId>
    <version>2.7.2version>
    <scope>providedscope>
dependency>

从socket中读取数据

import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

/**
 * 从端口读取数据
 */
public class Flink03_Source_Socket {
    public static void main(String[] args) throws Exception {
        //1、获取执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);//若 并行度1 按顺序读取数据,否则不按顺序
        //2、从端口读取数据
        DataStreamSource socketTextStream = env.socketTextStream("hadoop103", 9998);

        //3、打印
        socketTextStream.print();
        //4、执行
        env.execute();
    }
}

从kafka中读取数据

引入依赖

<dependency>
    <groupId>org.apache.flinkgroupId>
    <artifactId>flink-connector-kafka_2.11artifactId>
    <version>1.12.0version>
dependency>
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;

import java.util.Properties;

/**
 * kafka 读取数据
 */
public class Flink04_Source_Kafka {
    public static void main(String[] args) throws Exception {
        //1、获取执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);//若 并行度1 按顺序读取数据,否则不按顺序
        //2、配置 kafka
        Properties properties = new Properties();
        properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "hadoop103:9092");
        properties.put(ConsumerConfig.GROUP_ID_CONFIG, "Flink0214");
        //3、读取数据
        DataStreamSource test =
                env.addSource(new FlinkKafkaConsumer("test", new SimpleStringSchema(), properties));
        //4、打印
        test.print();
        env.execute();
        //主机上 测试 bin/kafka-console-producer.sh --topic test --broker-list hadoop103
    }
}

自定义source

import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.source.SourceFunction;
import org.wdh01.bean.WaterSensor;

import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
import java.net.Socket;
import java.nio.charset.StandardCharsets;

/**
 * 自定义 source
 */
public class Flink05_Source_MySource {
    public static void main(String[] args) throws Exception {
        //1、获取执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);//若 并行度1 按顺序读取数据,否则不按顺序
        //2、读取数据
        DataStreamSource ds103 = env.addSource(new MySource("hadoop103", 9998));
        //3、打印
        ds103.print();
        //4、执行
        env.execute();

    }

    /**
     * 自定义 Source,模拟 端口 source
     */
    public static class MySource implements SourceFunction {
        //定义属性 主机和 端口
        public String host;
        public Integer port;
        private boolean flag = true;
        Socket socket = null;
        BufferedReader reader = null;

        public MySource() {
        }

        public MySource(String host, Integer port) {
            this.host = host;
            this.port = port;
        }

        public void run(SourceContext ctx) throws Exception {
            //创建输入流
            socket = new Socket(host, port);
            reader = new
                    BufferedReader(new InputStreamReader(socket.getInputStream(), StandardCharsets.UTF_8));
            while (flag) {
                // 读取数据
                String line = reader.readLine();
                while (flag && line != null) {
                    //接收数据 并发送 flink
                    String[] s = line.split(",");
                    WaterSensor waterSensor = new WaterSensor(s[0], Long.parseLong(s[1]), Integer.parseInt(s[2]));
                    ctx.collect(waterSensor);
                    line = reader.readLine();
                }
            }
        }

        public void cancel() {
            flag = false;
            try {
                reader.close();
            } catch (IOException e) {
                e.printStackTrace();
            }
            try {
                socket.close();
            } catch (IOException e) {
                e.printStackTrace();
            }

        }
    }
}