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、准备集合 ListwaterSensors = 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、从文件读取数据 DataStreamSourcestringDataStreamSource = 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 文件系统读取数据,- 参数可以是目录也可以是文件
- 路径可以是相对路径也可以是绝对路径
- 相对路径是从系统属性user.dir获取路径: idea下是project的根目录, standalone模式下是集群节点根目录
- 也可以从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、从端口读取数据 DataStreamSourcesocketTextStream = 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、读取数据 DataStreamSourcetest = 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、读取数据 DataStreamSourceds103 = 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(); } } } }