Skip to main content

大数据-Flink

简介

Apache Flink是一个框架和分布式处理引擎,用于在无界和有界数据流上进行有状态计算。Flink被设计为在所有常见的集群环境中运行,以内存中的速度和任何规模执行计算。以下操作指导以1.12.0版本为例。

部署

sw_64架构暂不支持

197执行如下命令,安装flink包。

yum install flink -y

yum install java-1.8.0* -y

198执行如下命令,配置JAVA_HOME。

tail -n 4 /etc/profile

199执行如下命令,启动集群。

cd /opt/apache-flink-1.12.0

./bin/start-cluster.sh

访问web页面

打开浏览器,访问http://${localhost}:8081检查分派器的web前端,并确保一切正常,从而启动一个只有一个JobManager和TaskManager的本地Flink集群。

查看日志

您还可以执行如下命令,通过检查logs目录中的日志文件来验证系统是否正在运行:

tail log/flink--standalonesession-.log

运行实例

每个Flink的binary release都会包含一个examples(示例)目录,在其中可以找到这个页面上每个示例的jar包文件。

[root@localhost examples]# ls ./*

./batch:

ConnectedComponents.jar DistCp.jar EnumTriangles.jar KMeans.jar PageRank.jar TransitiveClosure.jar WebLogAnalysis.jar WordCount.jar

./gelly:

./python:

table

./streaming:

IncrementalLearning.jar SessionWindowing.jar StateMachineExample.jar Twitter.jar WordCount.jar

Iteration.jar SocketWindowWordCount.jar TopSpeedWindowing.jar WindowJoin.jar

./table:

ChangelogSocketExample.jar StreamSQLExample.jar StreamWindowSQLExample.jar WordCountSQL.jar

GettingStartedExample.jar StreamTableExample.jar TPCHQuery3Table.jar WordCountTable.jar

实例1 WordCount

WordCount是大数据系统中的“Hello World”,他可以计算一个文本集合中不同单词的出现频次。可以通过执行以下命令来运行WordCount示例:

:::note 其他示例也可通过类似方式执行。注意很多示例在不传递执行参数的情况下都会使用内置数据。 :::

如果需要利用WordCount程序计算真实数据,您需要传递存储数据的文件路径。例如:

./bin/flink run ./examples/batch/WordCount.jar --input ~/data --output ~/result

Job has been submitted with JobID e8630f0b63120d501b92d85e733e51d6

Program execution finished

Job with JobID e8630f0b63120d501b92d85e733e51d6 has finished.

Job Runtime: 462 ms

执行上述操作后,您可以检查Web界面以验证作业是否按预期运行:

您还可以执行如下命令,查看运行结果。

  • 输入文件

cat ~/data

hello

how are you

fine

thank you

and you

I am

fine too

  • 输出结果

cat ~/result

am 1

and 1

are 1

fine 2

hello 1

how 1

i 1

thank 1

too 1

you 3

实例2 SocketWindowWordCount

现在,我们将运行此Flink应用程序。它将从套接字读取文本,每5秒钟打印一次在前5秒钟内每个单词出现的次数,即只要单词漂浮在其中,处理时间就会翻腾。

200执行如下命令,启动本地服务器。

$ nc -l 9000

201重新打开一个终端,提交Flink程序:

Job has been submitted with JobID 2b2b0bf4ea55e45c17d0de52fb3a202b

202回到刚才的netcat界面,输入一些字符,例如:

lorem ipsum

ipsum ipsum ipsum

bye

203按“Ctrl+c”快捷键结束输入,可以看到第二个界面的程序执行成功。

[root@localhost apache-flink-1.12.0]# ./bin/flink run examples/streaming/SocketWindowWordCount.jar --port 9000

Job has been submitted with JobID 2b2b0bf4ea55e45c17d0de52fb3a202b

Program execution finished

Job with JobID 2b2b0bf4ea55e45c17d0de52fb3a202b has finished.

Job Runtime: 123983 ms

204检查Web界面,以验证作业是否按预期运行。

205执行如下命令,查看统计结果。

tail log/flink--taskexecutor-.out

lorem : 1

ipsum : 1

ipsum : 3

bye : 1

停止集群

./bin/stop-cluster.sh

Stopping taskexecutor daemon (pid: 18466) on host localhost.localdomain.

Stopping standalonesession daemon (pid: 18209) on host localhost.localdomain.