前言
今天更新 Flink 核心知识点!结合实际可运行的 3 个核心代码(EnvDemo、Event、SourceKafkaTest),提炼实操性极强的核心考点,没有晦涩理论,全是入门必懂的实操逻辑,适合 Flink 新手入门和代码复现练习,收藏起来跟着敲代码就能上手~
一、Flink 核心基础概念
1. 核心角色与执行环境
(1)执行环境(StreamExecutionEnvironment):Flink 程序的入口,负责初始化运行环境、接收作业、调度任务。
核心创建方式:StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();(自动适配本地 / 集群环境)
(2)关键配置:
①env.setParallelism(1):设置并行度(默认并行度为 CPU 核心数,开发环境建议设为 1,避免多线程干扰)
②env.setRuntimeMode(RuntimeExecutionMode.STREAMING):指定执行模式(STREAMING 流处理 / BATCH 批处理 / AUTO 自动识别)
(3)作业执行:Flink 程序为懒执行,必须调用env.execute("jobName")触发执行,参数为作业名称(如 "flinkJob")。
2. 核心数据载体:POJO 类(以 Event 为例)
Flink 中自定义数据类型需遵循 POJO 规范,确保序列化和反序列化正常:
(1)类为public修饰
(2)包含无参构造器(默认或显式定义)
(3)所有属性为public,或提供对应的getter/setter方法
(4)属性类型支持序列化(如 String、Long、基本类型等)
(5)示例核心结构:
public class Event {
public String user; // 用户名
public String url; // 访问地址
public Long timestamp; // 时间戳
public Event() {} // 无参构造器(必需)
public Event(String user, String url, Long timestamp) {
// 带参构造器(初始化用)
this.user = user;
this.url = url;
this.timestamp = timestamp;
}
// getter/setter方法(可选,属性public可省略)
// toString方法(可选,便于打印输出)
}
二、Flink 数据源(Source)实操
Source 是 Flink 程序的数据输入源,核心分为 4 类常用场景,对应 EnvDemo 和 SourceKafkaTest 中的实操案例:
1. 从集合读取(测试常用)
适用场景:本地测试、少量固定数据
核心 API:env.fromCollection(Collection<T> data)
示例代码:
// 初始化数据
Event tom = new Event("tom", "/home", 17989808080L);
Event jack = new Event("jack", "/cart", 17989808080L);
List<Event> dataList = new ArrayList<>();
dataList.add(tom);
dataList.add(jack);
// 从集合创建数据源
DataStreamSource<Event> source = env.fromCollection(dataList);
2. 从文件读取
适用场景:批处理、离线数据导入
核心 API:env.readTextFile(String filePath)(支持本地文件 / HDFS 文件路径)
示例代码:
// 读取本地input目录下的words.txt文件
DataStreamSource<String> fileSource = env.readTextFile("input/words.txt");
3. 从 Socket 读取(实时测试常用)
适用场景:本地实时测试,模拟无界数据流
核心 API:env.socketTextStream(String host, int port)(host 为服务端 IP,port 为端口号)
示例代码:
// 连接s2主机的6602端口读取数据
DataStreamSource<String> socketSource = env.socketTextStream("s2", 6602);
4. 从 Kafka 读取(生产环境核心)
(1)适用场景:生产环境实时数据流(高吞吐、高可用)
核心依赖(需在 pom.xml 中添加):
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka_${scala.binary.version}</artifactId>
<version>${flink.version}</version>
</dependency>
(2)核心步骤:
配置 Kafka 连接参数(bootstrap.servers、group.id 等)
创建 FlinkKafkaConsumer(指定 topic、反序列化器、配置)
通过env.addSource(consumer)添加数据源
示例代码(对应 SourceKafkaTest):
// 1. 配置Kafka参数
Properties props = new Properties();
props.setProperty("bootstrap.servers", "s1:9092,s2:9092,s3:9092"); // Kafka集群地址
props.setProperty("group.id", "consumer-group1"); // 消费者组ID
props.setProperty("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.setProperty("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
// 2. 创建Kafka消费者(topic为top_gas,反序列化器为SimpleStringSchema)
FlinkKafkaConsumer<String> kafkaConsumer = new FlinkKafkaConsumer<>
("top_gas", new SimpleStringSchema(), props);
// 3. 添加Kafka数据源
DataStreamSource<String> kafkaSource = env.addSource(kafkaConsumer);
三、数据输出(Sink)实操
Sink 是 Flink 程序的数据输出目的地,核心常用输出方式:
1. 打印输出(测试常用)
核心 API:dataStream.print()(默认打印到控制台,可加前缀区分来源)
示例代码:
socketSource.print(); // 直接打印
kafkaSource.print("Kafka"); // 加前缀"Kafka",输出格式:Kafka> 数据内容
2. 其他输出(生产环境扩展)
文件输出:StreamingFileSink(支持行编码 / 批量编码,适合批处理输出)
数据库输出:JDBC Sink(需添加 flink-connector-jdbc 依赖,写入 MySQL/PostgreSQL 等)
消息队列输出:Kafka Sink(将处理结果写回 Kafka,形成数据闭环)
四、核心配置与实操注意事项
1. 并行度配置
(1)优先级:算子级(operator.setParallelism(2))> 环境级(env.setParallelism(1))> 配置文件级(flink-conf.yaml 中的parallelism.default)
(2)开发建议:本地测试设为 1,避免多线程导致的输出乱序;生产环境根据集群资源调整(如 CPU 核心数、任务复杂度)。
2. 序列化注意事项
(1)自定义数据类型必须遵循 POJO 规范(如 Event 类),否则需手动指定序列化器(如 Kryo)
(2)避免使用无法序列化的类型(如 Thread、InputStream 等)作为 POJO 属性。
3. Kafka 连接关键配置
(1)bootstrap.servers:Kafka 集群的 broker 地址,多个地址用逗号分隔
(2)group.id:消费者组 ID,相同 ID 的消费者会分摊消费分区
(3)auto.offset.reset:消费偏移量重置策略(如latest从最新数据开始消费,earliest从最早数据开始消费)
(4)反序列化器:需与 Kafka 生产者的序列化器对应(如 StringDeserializer 对应 StringSerializer)
五、本部分核心实操考点
1. 执行环境的创建与配置(并行度、执行模式)
2. POJO 类的规范定义(无参构造器、属性序列化)
3. 三类核心 Source 的实现(集合 / 文件 / Socket)
4. Kafka Source 的配置与集成(生产环境重点)
5. Sink 打印输出的使用(测试排错必备)
结尾
这部分内容聚焦 Flink 入门最核心的「环境搭建 + 数据源 + 输出」实操,跟着代码敲一遍就能掌握 Flink 程序的基本骨架,为后续学习复杂算子和状态管理打下坚实基础。


