Flink_Source算子详解
Flink Source算子详解:从集合、文件、Socket到Kafka与数据生成器
前言
在Flink程序中,Source是整个数据处理流程的起点,负责从外部系统读取数据并构建DataStream。Flink提供了多种内置的Source实现,从简单的集合、文件,到Socket、Kafka等外部数据源,覆盖了开发测试和生产环境的各种场景。
本文基于尚硅谷Flink1.17教程,系统梳理Flink各种Source的用法与适用场景。
一、准备工作
本文使用WaterSensor(水位传感器)作为示例数据模型,字段如下:
| 字段名 | 数据类型 | 说明 |
|---|---|---|
| id | String | 传感器类型 |
| ts | Long | 记录时间戳 |
| vc | Integer | 水位值 |
public class WaterSensor {
public String id;
public Long ts;
public Integer vc;
public WaterSensor() {}
public WaterSensor(String id, Long ts, Integer vc) {
this.id = id;
this.ts = ts;
this.vc = vc;
}
@Override
public String toString() {
return "WaterSensor{id='" + id + "', ts=" + ts + ", vc=" + vc + '}';
}
}
这个类满足Flink POJO类型的要求:
- 类是
public的 - 有无参构造方法
- 所有属性是
public的 - 所有属性类型可序列化
Flink会把满足这些条件的类当作POJO来处理,能自动生成高效的序列化器。
二、新旧Source API的区别
在开始之前,先了解一个重要变化:
Flink 1.12之前,添加Source使用addSource()方法:
DataStream<String> stream = env.addSource(...);
// 传入的是SourceFunction接口的实现
Flink 1.12之后,推荐使用流批统一的新Source架构:
DataStreamSource<String> stream = env.fromSource(...);
新API统一了流批处理的Source接口,后续示例均使用新API。
三、从集合中读取数据
用法
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
List<Integer> data = Arrays.asList(1, 22, 3);
DataStreamSource<Integer> ds = env.fromCollection(data);
ds.print();
env.execute();
也可以直接使用fromElements读取零散元素:
DataStreamSource<WaterSensor> stream = env.fromElements(
new WaterSensor("sensor_1", 1L, 1),
new WaterSensor("sensor_2", 2L, 2),
new WaterSensor("sensor_3", 3L, 3)
);
适用场景
仅用于开发测试,数据直接存在内存中,方便验证逻辑是否正确,不适合生产环境。
四、从文件读取数据
添加依赖
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-files</artifactId>
<version>${flink.version}</version>
</dependency>
用法
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
FileSource<String> fileSource = FileSource
.forRecordStreamFormat(new TextLineInputFormat(), new Path("input/word.txt"))
.build();
env.fromSource(fileSource, WatermarkStrategy.noWatermarks(), "file-source")
.print();
env.execute();
几个注意点
- 路径参数可以是文件,也可以是目录(会读取目录下所有文件)
- 支持HDFS路径,格式为
hdfs://hadoop102:9000/path/to/file - 相对路径的基准目录:IDEA中是项目根目录,Standalone集群中是节点根目录
适用场景
适合批处理场景,读取日志文件、历史数据等有界数据。
五、从Socket读取数据
用法
DataStream<String> stream = env.socketTextStream("localhost", 7777);
启动后在终端使用nc -lk 7777发送数据,Flink就能实时接收。
适用场景
读取的是无界数据流,是真正的流处理场景。但由于吞吐量小、稳定性差,只适合开发测试,不能用于生产环境。
六、从Kafka读取数据
Kafka是生产环境中最常用的数据源,Flink官方提供了flink-connector-kafka连接器。
添加依赖
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>${flink.version}</version>
</dependency>
用法
public class SourceKafka {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
.setBootstrapServers("hadoop102:9092") // Kafka地址
.setTopics("topic_1") // 消费的Topic
.setGroupId("atguigu") // 消费者组
.setStartingOffsets(OffsetsInitializer.latest()) // 从最新offset开始消费
.setValueOnlyDeserializer(new SimpleStringSchema()) // 反序列化方式
.build();
DataStreamSource<String> stream = env.fromSource(
kafkaSource,
WatermarkStrategy.noWatermarks(),
"kafka-source"
);
stream.print("Kafka");
env.execute();
}
}
核心参数说明
| 参数 | 说明 |
|---|---|
setBootstrapServers |
Kafka集群地址,多个用逗号分隔 |
setTopics |
消费的Topic名称,可以传多个 |
setGroupId |
消费者组ID |
setStartingOffsets |
起始消费位置(见下表) |
setValueOnlyDeserializer |
只对value进行反序列化 |
起始offset常用选项:
| 方法 | 含义 |
|---|---|
OffsetsInitializer.latest() |
从最新消息开始消费 |
OffsetsInitializer.earliest() |
从最早消息开始消费 |
OffsetsInitializer.committedOffsets() |
从已提交的offset继续消费 |
适用场景
生产环境中最主流的数据源,适合实时流处理场景。关于Kafka的深入使用(消费者组原理、offset管理、水位线配合等),后续会单独写一篇详细介绍。
七、从数据生成器读取数据
Flink从1.11开始提供了内置的DataGen连接器,用于生成随机数据,主要用于没有数据源时的性能测试和流任务验证。
添加依赖
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-datagen</artifactId>
<version>${flink.version}</version>
</dependency>
用法
public class DataGeneratorDemo {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
DataGeneratorSource<String> dataGeneratorSource = new DataGeneratorSource<>(
new GeneratorFunction<Long, String>() {
@Override
public String map(Long value) throws Exception {
return "Number:" + value; // 自定义生成逻辑
}
},
Long.MAX_VALUE, // 生成数据的总条数上限
RateLimiterStrategy.perSecond(10), // 每秒生成10条
Types.STRING // 输出数据类型
);
env.fromSource(dataGeneratorSource, WatermarkStrategy.noWatermarks(), "datagenerator")
.print();
env.execute();
}
}
核心参数说明
| 参数 | 说明 |
|---|---|
GeneratorFunction |
定义数据生成逻辑,入参是自增的Long序号 |
| 第二个参数(count) | 生成数据的总条数,Long.MAX_VALUE表示无限生成 |
RateLimiterStrategy |
限流策略,perSecond(n)表示每秒生成n条 |
Types.STRING |
输出数据的类型信息 |
适用场景
无数据源时的测试和性能压测,可以精确控制数据生成速率和总量。
八、Flink支持的数据类型
使用Source读取数据时,需要了解Flink支持哪些数据类型。
支持的类型
Flink使用TypeInformation来统一表示数据类型,支持以下几类:
| 类型 | 说明 |
|---|---|
| 基本类型 | Java基本类型及包装类,String、Date、BigDecimal等 |
| 数组类型 | 基本类型数组、对象数组 |
| Tuple类型 | Flink内置元组,支持Tuple0~Tuple25 |
| ROW类型 | 任意字段数的元组,支持空字段 |
| POJO类型 | 满足特定条件的Java Bean |
| 辅助类型 | Option、Either、List、Map等 |
| 泛型类型 | 不满足上述条件的类,由Kryo序列化 |
推荐使用POJO类型,支持按字段名定义key,代码可读性最好。
类型提示(Type Hints)
由于Java泛型擦除,某些情况下(尤其是Lambda表达式)Flink无法自动推断类型,需要手动指定:
// 使用.returns()指定返回类型
.map(word -> Tuple2.of(word, 1L))
.returns(Types.TUPLE(Types.STRING, Types.LONG));
// 或者使用TypeHint
.returns(new TypeHint<Tuple2<String, Long>>(){});
九、各Source对比总结
| Source类型 | 数据特点 | 适用场景 | 是否生产可用 |
|---|---|---|---|
| 集合/fromElements | 有界、内存 | 单元测试、逻辑验证 | ❌ |
| 文件 | 有界 | 批处理、历史数据分析 | ✅(批处理) |
| Socket | 无界、低吞吐 | 开发测试 | ❌ |
| Kafka | 无界、高吞吐 | 实时流处理 | ✅ |
| 数据生成器 | 无界、可控速率 | 性能压测、无数据源测试 | ❌ |
小结
- 测试用:集合、Socket、数据生成器,根据是否需要控制速率和数据量选择
- 批处理用:文件Source,支持本地文件和HDFS
- 生产流处理用:Kafka Source,高吞吐、可靠、支持offset管理
理解各种Source的特点和适用场景,是写Flink程序的第一步。后续文章会深入介绍Kafka Source的进阶用法,包括offset管理策略、水位线配合以及反序列化的自定义实现。
更多推荐
所有评论(0)