什么是Flink

Flink 是一个分布式流处理引擎,专为处理无界(实时)和有界(批处理)数据流设计,以低延迟、高吞吐和精准的 Exactly-Once 语义著称。

它的核心特点是流批一体—— 用同一套引擎同时处理实时流数据和历史批数据,无需切换框架,且内置强大的状态管理、事件时间语义和容错机制,能保证数据处理的准确性和一致性。

典型应用场景包括实时风控(如信用卡欺诈检测)、实时报表(如电商大屏)、实时推荐(如短视频内容推送)等。

Flink 安装

Flink 旨在以极快的速度处理连续数据流。本简短指南将向你展示如何下载、安装并运行 Flink 的最新稳定版本。你还将运行一个 Flink 示例作业,并在 Web 界面中查看它。

下载Flink

Flink 可在所有类 UNIX 环境中运行,即 Linux、Mac OS X 以及适用于 Windows 的 Cygwin。你需要安装 Java 11。要检查已安装的 Java 版本,请在终端中输入:

yangyanping@yangyaningdeAir flink-2.2.0 % java -version
java version "17.0.15" 2025-04-15 LTS
Java(TM) SE Runtime Environment (build 17.0.15+9-LTS-241)
Java HotSpot(TM) 64-Bit Server VM (build 17.0.15+9-LTS-241, mixed mode, sharing)

接下来,下载 Flink 的最新二进制发行版https://flink.apache.org/downloads/,然后解压归档文件:

yangyanping@yangyaningdeAir Downloads % tar -xzf flink-2.2.0-bin-scala_2.12.tgz 
yangyanping@yangyaningdeAir Downloads % cd flink-2.2.0
yangyanping@yangyaningdeAir flink-2.2.0 % ls
LICENSE		README.txt	conf		lib		log		plugins
NOTICE		bin		examples	licenses	opt

浏览项目目录 

进入解压后的目录,并通过以下命令列出内容:

yangyanping@yangyaningdeAir Downloads % cd flink-2.2.0 && ls -l
total 376
-rw-r--r--@  1 yangyanping  staff   11357 11 10 10:18 LICENSE
-rw-r--r--@  1 yangyanping  staff  172202 11 27 14:52 NOTICE
-rw-r--r--@  1 yangyanping  staff    1309 11 10 10:18 README.txt
drwxr-xr-x@ 27 yangyanping  staff     864 11 27 14:52 bin
drwxr-xr-x@ 13 yangyanping  staff     416 11 27 14:52 conf
drwxr-xr-x@  5 yangyanping  staff     160 11 27 14:52 examples
drwxr-xr-x@ 16 yangyanping  staff     512 11 27 14:52 lib
drwxr-xr-x@ 38 yangyanping  staff    1216 11 27 14:52 licenses
drwxr-xr-x@  8 yangyanping  staff     256  2 26 20:17 log
drwxr-xr-x@ 15 yangyanping  staff     480 11 27 14:52 opt
drwxr-xr-x@ 12 yangyanping  staff     384 11 27 14:52 plugins

现在,你可能需要注意以下几点:

  • bin/ 目录包含 Flink 二进制文件以及多个用于管理各种作业和任务的 bash 脚本
  • conf/ 目录包含配置文件,其中包括 Flink 配置文件
  • examples/ 目录包含可直接与 Flink 一起使用的示例应用程序

启动和停止本地集群

要启动本地集群,请运行 Flink 附带的 bash 脚本:

yangyanping@yangyaningdeAir bin % ./start-cluster.sh 
Starting cluster.
Starting standalonesession daemon on host yangyaningdeAir.
Starting taskexecutor daemon on host yangyaningdeAir.

Flink 现在正在后台进程中运行。你可以使用以下命令检查其状态:

$ ps aux | grep flink

你可以访问位于 localhost:8081 的 Web 界面,查看 Flink 仪表板,确认集群已启动并正常运行。

若要快速停止集群及所有运行中的组件,可使用提供的脚本:

$ ./bin/stop-cluster.sh

提交 Flink 作业

Flink 提供了一个命令行工具 bin/flink,可运行打包为 Java 归档文件(JAR)的程序并控制其执行。提交作业意味着将作业的 JAR 文件及相关依赖项上传到正在运行的 Flink 集群并执行它。

Flink 发行版附带示例作业,你可以在 examples/ 文件夹中找到它们。

要将示例单词计数作业部署到正在运行的集群,请执行以下命令:

yangyanping@yangyaningdeAir flink-2.2.0 % ./bin/flink run examples/streaming/WordCount.jar 
Executing example with default input data.
Use --input to specify file input.
Printing result to stdout. Use --output to specify output path.
Job has been submitted with JobID 0368312bb6fc13cf1f10db129bf93cb5
Program execution finished
Job with JobID 0368312bb6fc13cf1f10db129bf93cb5 has finished.
Job Runtime: 579 ms

你可以通过查看日志来验证输出:(默认日志位于 Flink 解压目录下的 logs/ 文件夹中)

yangyanping@yangyaningdeAir flink-2.2.0 % tail log/flink-yangyanping-taskexecutor-0-yangyaningdeAir.out 
(nymph,1)
(in,3)
(thy,1)
(orisons,1)
(be,4)
(all,2)
(my,1)
(sins,1)
(remember,1)
(d,4)

此外,你可以访问 Flink 的 Web UI,监控集群状态与运行中作业的执行情况。

你还可以查看本次作业执行对应的数据流执行计划:

在本次作业的执行流程中,Flink 定义了两个算子。第一个是源算子,负责从集合数据源中读取数据;第二个是转换算子,负责对单词计数结果进行聚合统计。想要了解更多细节,可查阅 DataStream 算子的相关文档。

你还可以查看该作业执行的时间线:

你已成功运行了一个 Flink 应用程序!你可以自由选用 examples/ 文件夹下的其他 JAR 归档文件,或是部署你自己开发的作业!

总结

在本指南中,你完成了 Flink 的下载、项目目录的浏览、本地集群的启动与停止,还成功提交了一个 Flink 示例作业!

更多推荐