
(1)文件数据源
在 StreamExecutionEnvironment 中,可以使用 readTextFile 方法直接读取文本文件,也可以使用 readFile 方法通过指定文件 InputFormat 来读取特定数据类型的文件,如 CsvInputFormat。
下面的代码演示了使用 readTextFile 读取文本文件
import org.apache.flink.streaming.api.scala._object Flink9 extends App {val env = StreamExecutionEnvironment.getExecutionEnvironmentenv.readTextFile("d://1.txt").print()env.execute("read_textfile")}
(2)Socket 数据源
env.socketTextStream("localhost",9999)
在 unix 环境下,可以执行 nc -lk 9999 命令,启动端口,在客户端中输入数据,flink 就能接收到数据了
(3)集合数据源
可以直接将 Java 或 Scala 程序中的集合类 转换成 DataStream 数据集,本质上是将本地集合中的数据分发到远端并行执行的节点中。如:
val dataStream = env.fromElements(Tuple2(1L,3L),Tuple2(1L,5L))val dataStream2 = env.fromCollection(List(1,2,3))

2
外部数据源
前面的数据源类型都是非常基础的数据接入方式,例如从文件,Socket 端口中接入数据,其本质是实现了不同的 SourceFunction,Flink 将其封装成高级的 API,减少了用户的使用成本。
企业中,大部分都是使用高性能的第三方存储介质和中间件,比如 Kafka,Elasticsearch,RabbitMQ 等。
下面以 Kafka 为例,来说明如何使用 kafka 作为 输入源。
首先,需要在 pom.xml 文件中引入依赖
<dependency><groupId>org.apache.flink</groupId><artifactId>flink-connector-kafka-0.10_2.11</artifactId><version>1.8.0</version></dependency>
(我这里使用的是 kafka 0.10 ,flink 1.8.0 版本的)
引入 maven 配置后,就可以在 Flink 应用工程中创建和使用相应的 Connector了,主要的参数有 kafka topic,bootstrap.servers,zookeeper.connect,另外 Schema 参数的主要作用是根据事先定义好的 Schema 信息将数据序列化成该 Schema 定义的数据类型,默认是 SimpleStreamSchema,代表从 Kafka 中接入的数据转换成 String 类型。
package com.dsj361import java.util.Propertiesimport org.apache.flink.api.common.serialization.SimpleStringSchemaimport org.apache.flink.streaming.api.scala._import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer010object Flink10 extends App {val env = StreamExecutionEnvironment.getExecutionEnvironmentval properties = new Properties()properties.setProperty("bootstrap.servers","192.168.17.21:9092")properties.setProperty("zookeeper.connect","192.168.17.22:2181")properties.setProperty("group.id","test")val input =env.addSource(new FlinkKafkaConsumer010[String]("topic1",new SimpleStringSchema(),properties))input.print()env.execute("connect_kafka")}
同时,我们可以把 从 kafka 读到的数据反序列化成我们期望的格式,主要是实现 DeserializationSchema 来完成。
Flink 中已经实现了大多数主流的数据源连接器,但是 Flink 的整体架构非常开放,用户可以自定义连接器,以满足不同数据源的接入需求。
可以通过实现 SourceFunction 定义单个线程的数据接入器,也可以通过实现 ParallelSourceFunction 接口 或者继承 RichParallelSourceFunction 类定义并发数据源接入器
(关于 kafka 的接入会单独开辟一张来讲解)
下一讲预告:DataStream 的各种花式转换操作。

为了方便大家共同交流学习
建立了 QQ 群(780016495),大家可以扫码进入,互相探讨交流

另外关注个人公众号,获得持续的 Flink 技术更新,扫码关注即可





