大数据-Flink
大数据-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:
flink-gelly-examples_2.11-1.12.0.jar
./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示例:
./bin/flink run ./examples/batch/WordCount.jar
:::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程序:
$ ./bin/flink run examples/streaming/SocketWindowWordCount.jar --port 9000
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.
No comments to display
No comments to display