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管理策略、水位线配合以及反序列化的自定义实现。

更多推荐