欢迎光临
我们一直在努力

Flink 核心知识点

前言

今天更新 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 程序的基本骨架,为后续学习复杂算子和状态管理打下坚实基础。

赞(0)
未经允许不得转载:171主机测评 » Flink 核心知识点
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址