从零开始学Flink:数据源

📅 2026/7/25 18:11:58
从零开始学Flink:数据源
从零开始学Flink数据源前言Apache Flink 是一个分布式流处理框架以其高吞吐、低延迟和精确一次语义Exactly-Once在实时计算领域占据重要地位。数据源Source是 Flink 作业的起点——它定义了数据从哪里来。本文将从实战出发带你一步步理解 Flink 数据源的原理、分类并通过完整代码示例掌握常见数据源的用法。## 1. Flink 数据源分类Flink 的数据源主要分为三类-内置数据源如文件、Socket、集合等适合测试和快速开发。-自定义数据源通过实现SourceFunction或RichSourceFunction接口满足业务特定需求。-连接器数据源如 Kafka、RabbitMQ 等与外部系统集成。在实际项目中Kafka 是最常用的数据源但为了从零开始学习我们先从简单场景入手。## 2. 搭建开发环境本文代码基于 Flink 1.14 和 Java 11。假设你已经安装 Maven 和 JDK。在pom.xml中添加 Flink 依赖xmldependencies dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java_2.12/artifactId version1.14.0/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-clients_2.12/artifactId version1.14.0/version /dependency/dependencies## 3. 第一个数据源从集合读取集合数据源是最简单的入门方式适合单元测试或原型开发。javaimport org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;public class CollectionSourceExample { public static void main(String[] args) throws Exception { // 1. 创建执行环境 StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 2. 从集合创建数据源 DataStreamString textStream env.fromElements( hello flink, hello world, flink streaming ); // 3. 转换操作每行按空格切分并统计词频 DataStreamString wordCounts textStream .flatMap((String line, org.apache.flink.util.CollectorString out) - { for (String word : line.split( )) { out.collect(word); } }) .returns(String.class) .map(word - word count: 1); // 4. 打印结果并行度设为1以便观察 wordCounts.print().setParallelism(1); // 5. 执行作业 env.execute(Collection Source Example); }}运行结果控制台输出hello count: 1flink count: 1hello count: 1world count: 1flink count: 1streaming count: 1关键点fromElements方法直接接收可变参数底层通过CollectionSource将集合转换为数据流。此方法适用于数据量小且无状态转换的场景。## 4. 第二个数据源从 Socket 读取Socket 数据源常用于模拟实时流数据。你可以通过nc -lk 9999命令启动一个端口然后手动输入数据。javaimport org.apache.flink.api.common.functions.FlatMapFunction;import org.apache.flink.api.java.tuple.Tuple2;import org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.util.Collector;public class SocketSourceExample { public static void main(String[] args) throws Exception { // 1. 创建执行环境 StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 2. 从Socket读取数据需先启动 nc -lk 9999 DataStreamString socketStream env.socketTextStream(localhost, 9999); // 3. 实时词频统计 DataStreamTuple2String, Integer wordCounts socketStream .flatMap(new FlatMapFunctionString, Tuple2String, Integer() { Override public void flatMap(String line, CollectorTuple2String, Integer out) throws Exception { for (String word : line.split(\\s)) { out.collect(new Tuple2(word, 1)); } } }) .keyBy(value - value.f0) // 按单词分组 .sum(1); // 对计数求和 // 4. 打印结果 wordCounts.print().setParallelism(1); // 5. 执行作业Socket数据源会一直等待输入 env.execute(Socket Source Example); }}运行步骤1. 在终端执行nc -lk 9999启动服务器。2. 运行 Java 程序。3. 在终端输入hello flink控制台输出(hello,1)和(flink,1)。4. 再次输入hello world输出(hello,2)和(world,1)。关键点socketTextStream返回一个DataStreamString数据是无限流除非连接断开。此示例展示了有状态的流处理——keyBy和sum共同实现了单词的累加统计。## 5. 自定义数据源模拟随机数据当内置数据源无法满足需求时可以自定义实现SourceFunction。下面是一个每1秒生成一条随机传感器数据的示例javaimport org.apache.flink.streaming.api.functions.source.SourceFunction;import org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import java.util.Random;public class CustomSourceExample { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 使用自定义数据源 DataStreamSensorReading sensorStream env.addSource(new SensorSource()); sensorStream.print().setParallelism(1); env.execute(Custom Source Example); } // 定义传感器读数类 public static class SensorReading { public String id; public long timestamp; public double temperature; public SensorReading(String id, long timestamp, double temperature) { this.id id; this.timestamp timestamp; this.temperature temperature; } Override public String toString() { return SensorReading{ id id \ , timestamp timestamp , temperature temperature }; } } // 自定义 SourceFunction public static class SensorSource implements SourceFunctionSensorReading { private volatile boolean running true; // 控制停止 private final Random random new Random(); Override public void run(SourceContextSensorReading ctx) throws Exception { while (running) { // 生成10个传感器的随机数据 for (int i 1; i 10; i) { String sensorId sensor_ i; long timestamp System.currentTimeMillis(); double temperature 20.0 random.nextGaussian() * 10; // 均值20标准差10 ctx.collect(new SensorReading(sensorId, timestamp, temperature)); } // 每隔1秒生成一批数据 Thread.sleep(1000); } } Override public void cancel() { running false; } }}运行结果示例输出SensorReading{idsensor_1, timestamp1680000000123, temperature22.345}SensorReading{idsensor_2, timestamp1680000000123, temperature18.789}...关键点- 实现SourceFunction需要重写run和cancel方法。-run方法中通过SourceContext.collect发射数据。-cancel方法用于优雅停止数据源例如作业取消时。- 通过volatile变量保证多线程可见性。## 6. 连接器数据源Kafka入门介绍虽然本文侧重基础但必须提及 Kafka 连接器。在pom.xml中添加xmldependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka_2.12/artifactId version1.14.0/version/dependency然后通过FlinkKafkaConsumer读取javaProperties props new Properties();props.setProperty(bootstrap.servers, localhost:9092);props.setProperty(group.id, flink-consumer);DataStreamString kafkaStream env.addSource( new FlinkKafkaConsumer(topic-name, new SimpleStringSchema(), props));Kafka 连接器支持精确一次语义、动态分区发现等高级特性是生产环境的首选。## 7. 数据源的并行度与性能- 内置数据源如集合、文件的并行度通常为1因为数据是有限的。- Socket 数据源并行度只能为1因为一个端口只能被一个线程监听。- 自定义数据源可以通过实现ParallelSourceFunction或RichParallelSourceFunction来支持多并行度。- 对于 Kafka 等高吞吐连接器适当提高并行度可以提升吞吐量。## 8. 常见问题与调试技巧1.数据源不输出数据检查是否调用了env.execute()。2.Socket 连接失败确保先启动nc命令。3.自定义数据源阻塞在run方法中不要使用while(true)而不加休眠或条件。4.并行度设置通过setParallelism(n)控制但某些数据源会忽略此设置。## 总结本文从零开始通过三个完整代码示例集合、Socket、自定义数据源演示了 Flink 数据源的基本用法。核心要点包括-内置数据源适合快速原型但生产环境需使用连接器。-自定义数据源通过实现SourceFunction提供灵活性可模拟任何数据输入。-Socket 数据源是学习流处理实时性的最佳入门选择。-Kafka 连接器是业界标准应作为后续学习的重点。数据源是 Flink 作业的入口掌握了它你就迈出了实时计算的第一步。建议读者动手运行上述代码并尝试修改参数如并行度、数据频率观察行为变化。下一篇文章将深入讨论数据转换算子敬请期待