从零上手Flink:这篇口语化教程,让你少走3小时弯路
想入门Flink,但看官方文档像看天书,咋办?其实我刚开始学的时候也一样——满屏的“流处理”“状态管理”,越看越懵。后来摸透了门道才发现,Flink入门没那么复杂,关键是把“为什么用”“怎么搭环境”“怎么写第一行代码”这几件事说清楚。今天就用唠嗑的方式,带你一步步走进Flink的世界,看完就能动手实操!
一、先搞懂:Flink到底是干啥的?为啥现在都用它?
咱先别着急啃技术细节,先搞明白Flink的“定位”。简单说,Flink是个专门处理“流数据”的框架——比如你手机里实时刷新的外卖订单、抖音的实时推荐、银行的转账流水,这些“一刻不停产生的数据”,都得靠Flink这种工具来处理。
可能有人会问:“我之前听过Spark Streaming,这不也是处理流数据的吗?为啥非要学Flink?”这里给大家掰扯两个核心区别,你就懂了:
1. 处理数据的“实时性”不一样:Spark Streaming是“伪实时”——它会把数据攒成一小批再处理,比如每隔1秒处理一次;而Flink是“真实时”——数据一来就立刻处理,延迟能做到毫秒级(简单理解:比Spark Streaming快100倍不止)。现在做实时推荐、实时风控,对延迟要求特别高,Flink自然就成了首选。
2. 处理数据的“完整性”不一样:比如统计“今天的订单总数”,如果中间服务器重启了,Spark Streaming可能会丢数据;但Flink有个“状态管理”的本事——能记住之前处理到哪了,重启后接着来,数据不丢也不重复。
总结下:如果你的工作涉及“实时数据处理”(比如实时报表、实时预警、实时推荐),那Flink几乎是绕不开的工具。现在大厂招聘里,“会Flink”已经成了大数据开发的加分项,这也是咱要学它的理由~
二、环境搭建:3步搞定,比装QQ还简单
很多人入门卡第一步就是“环境搭不好”,其实只要跟着步骤来,10分钟就能搞定。我用的是Windows系统,Mac步骤差不多,大家照着做就行。
1. 先装“前置条件”:JDK和Hadoop(可选)
Flink是用Java写的,所以必须先装JDK,而且得是JDK 8或11(别装太高版本,容易出兼容问题)。装完后记得配置环境变量(电脑右键→属性→高级系统设置→环境变量,加个JAVA_HOME,指向JDK安装目录),然后打开cmd输“java -version”,能显示版本号就成。
至于Hadoop,如果你只是本地练手,可以不装——Flink支持“本地模式”,直接跑就行;如果想模拟集群环境,再装Hadoop也不迟,这里先讲最简单的本地模式。
2. 下载Flink安装包:选对版本很重要
直接去Flink官网(https://flink.apache.org/)下载,注意两点:
- 版本别选最新的“Snapshot版”(开发中的版本,不稳定),选“Stable版”,比如我用的是Flink 1.17.0(2024年比较稳定的版本)。
- 下载“Scala 2.12”版本的(Scala是Flink的另一种开发语言,咱用Java的话,选2.12版本兼容性最好)。
下载完后解压到随便一个文件夹(别放有中文或空格的路径,比如“D:\flink-1.17.0”,不然容易报错)。
3. 启动Flink:双击图标就行,不用敲复杂命令
打开解压后的Flink文件夹,找到“bin”目录,里面有个“start-cluster.bat”(Windows)或“start-cluster.sh”(Mac/Linux),双击它!
然后打开浏览器,输入“http://localhost:8081”,如果能看到Flink的管理界面(左边有“Cluster Overview”,显示“Task Managers: 1”),说明环境搭好了!是不是比装QQ还简单?
三、实战:写第一个Flink程序,统计实时单词
环境搭好后,咱来写个最经典的“实时单词统计”程序——模拟从一个“数据流”里读单词,然后实时算出每个单词出现的次数。用IntelliJ IDEA来写,步骤很清晰。
1. 新建Maven项目,导入Flink依赖
打开IDEA,新建一个Maven项目(不用选模板,直接Next),然后在“pom.xml”里加Flink的依赖(复制粘贴就行,注意版本要和你下载的Flink一致):
xml
<<dependencies>
<!-- Flink核心依赖 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>1.17.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.17.0</version>
</dependency>
<!-- Flink本地执行依赖(本地练手必须加) -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients</artifactId>
<version>1.17.0</version>
</dependency>
</</dependencies>
加完后点击右上角的“刷新”按钮,让IDEA下载依赖(第一次可能有点慢,耐心等会)。
2. 写代码:50行搞定实时单词统计
新建一个Java类(比如叫“WordCountStreaming”),代码里每一步我都加了注释,大家跟着看就行:
java
import 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 WordCountStreaming {
public static void main(String[] args) throws Exception {
// 1. 获取Flink的执行环境(相当于“启动一个Flink实例”)
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 本地练手可以设置并行度为1(简单理解:只用一个线程处理,结果更直观)
env.setParallelism(1);
// 2. 定义“数据源”:这里用本地端口9999作为输入(等会我们用cmd发数据)
DataStream<String> inputDataStream = env.socketTextStream("localhost", 9999);
// 3. 处理数据:切分单词→转换成(单词,1)的格式→按单词分组→统计次数
DataStream<Tuple2<String, Integer>> resultStream = inputDataStream
// 切分单词:把每一行数据切成单个单词,比如“hello world”→“hello”“world”
.flatMap(new FlatMapFunction<String, Tuple2<String, Integer>>() {
@Override
public void flatMap(String line, Collector<Tuple2<String, Integer>> out) throws Exception {
// 按空格切分单词
String[] words = line.split(" ");
// 把每个单词转换成(单词,1),传给下一级
for (String word : words) {
out.collect(new Tuple2<>(word, 1));
}
}
})
// 按单词分组(Tuple2的第一个元素是单词,所以用0表示)
.keyBy(tuple -> tuple.f0)
// 统计次数(Tuple2的第二个元素是1,所以用1表示求和)
.sum(1);
// 4. 输出结果:把统计结果打印到控制台
resultStream.print();
// 5. 启动Flink程序(这行必须加,不然程序不会执行)
env.execute("Flink WordCount Streaming");
}
}
3. 运行程序,看实时效果
步骤分两步,很关键,别漏掉:
1. 启动本地端口监听:打开cmd,输入“nc -l -p 9999”(如果提示“nc不是内部命令”,说明没装netcat,去网上搜“Windows netcat下载”,装完后再试)。这个端口的作用是:我们在cmd里输入单词,Flink程序会实时读取。
2. 运行IDEA里的程序:点击main方法旁边的“运行”按钮,此时程序会处于“等待输入”的状态。
然后在刚才的cmd窗口里输入单词,比如“hello flink”“hello world”,再看IDEA的控制台——会实时显示每个单词的统计次数!比如输入“hello flink”后,控制台会输出(hello,1)(flink,1);再输入“hello world”,会输出(hello,2)(world,1)。这就是Flink的实时处理能力,是不是很直观?
四、入门必懂:3个核心概念,别被术语吓到
刚才写程序的时候,其实已经用到了Flink的核心概念,只是没细说。这里用大白话解释下,以后看文档就不懵了:
1. DataStream:你可以理解成“流动的数据管道”
比如刚才的“inputDataStream”,就是从端口9999读进来的“数据流”;“resultStream”就是处理完后的“结果数据流”。Flink的所有操作,都是在“数据流”上做的——就像水在管道里流,我们在管道中间加各种“处理器”(切分单词、分组、求和)。
2. 状态(State):Flink的“记忆力”
刚才统计单词次数的时候,Flink需要记住“hello”之前已经出现过1次,所以再收到“hello”时,才会变成2次。这个“记住之前的次数”的能力,就是“状态”。简单说,状态就是Flink保存的“中间计算结果”,就算程序重启,也能恢复,数据不丢。
3. 并行度(Parallelism):Flink的“多线程处理”
刚才在代码里加了“env.setParallelism(1)”,意思是用1个线程处理。如果把并行度改成2,Flink就会用2个线程同时处理数据,速度更快。比如处理100万条数据,并行度2比并行度1快一倍(前提是电脑CPU够)。
五、新手常见坑:我踩过的3个雷,你别再犯
入门的时候,我踩过几个坑,浪费了不少时间,这里提醒下大家:
1. JDK版本不对:比如装了JDK 17,Flink 1.17.0不兼容,会报“类找不到”的错。记住,JDK 8或11最稳妥。
2. 没启动nc端口:运行程序后,控制台没反应,一看才发现没启动“nc -l -p 9999”。Flink程序等着读数据,没数据源自然不干活。
3. 忘记写env.execute():代码里所有步骤都对,但程序就是不运行,最后发现少了“env.execute()”。这行代码是“启动开关”,必须加!
六、下一步:学会这些,你就超过80%的新手
入门之后,想进一步学习,可以从这3个方向入手:
1. 学“窗口(Window)”:比如“统计每5秒内的单词次数”,这就需要窗口——Flink里最常用的功能之一,用来处理“一段时间内的数据”。
2. 学“状态管理”:比如怎么保存状态、怎么恢复状态,这是Flink保证数据不丢的核心,面试常问。
3. 尝试“集群模式”:本地模式只能练手,实际工作中都是用集群(多台机器一起跑),可以学怎么把Flink程序提交到集群运行。
今天的Flink入门就到这了,其实核心就是“搭环境→写简单程序→理解核心概念”。刚开始练的时候,别追求一次性看懂所有细节,先跑通程序,再慢慢琢磨每个步骤的作用
最后说一句:Flink入门不难,难的是坚持练。把今天的单词统计程序改改,比如换成从文件读数据,或者统计每个用户的订单数,多动手,很快就能上手!
更多推荐


所有评论(0)