简介:本资源为大数据课程实验8「Flink初级编程实践」的完整实验报告,面向正在学习大数据技术原理与应用的高校学生及Flink入门开发者,帮助解决从环境搭建到作业提交的全流程实践问题。压缩包内为1个docx文档,约2.46MB,完整记录了实验环境配置、代码编写与运行结果。报告围绕两个核心任务展开:一是使用IntelliJ IDEA开发WordCount程序,涵盖Flink与Maven安装、Java代码编写、JAR打包及集群提交运行;二是借助Linux自带NC程序模拟实时数据流,编写Flink程序完成词频统计并部署运行。文中还整理了Idea引用Flink报错、Maven打包缓慢、NC程序无输出等常见问题的排查与解决思路,并附有Flink Web控制台查看输出的方法。目前已有5153人学习下载,适合需要对照实操、复盘流程与查漏补缺的初学者参考。
1. 从一份 2022 年的 Flink 实验报告说起:它到底能帮你省下多少试错时间
如果你正在搜「Flink 初级编程实践」,大概率是三种人之一:正在交大数据实验报告的学生、刚转流处理方向想跑通第一个 Demo 的工程师、或者被公司要求做实时词频统计原型的技术负责人。这份实验报告的核心价值不在于代码有多复杂,而在于它把 Flink 从安装到提交 JAR 包再到 NC 数据流模拟的完整链路走了一遍,而且记录了三个真实翻车现场:IDEA 引用 Flink 报错、Maven 打包慢到怀疑人生、NC 程序跑起来没输出。这三个坑几乎每个 Flink 新手都会踩,区别只是有人花了两小时,有人花了两天。实验环境是 Windows 10 加 VirtualBox 跑 Ubuntu 64-bit,内存给了 2048MB,处理器 4 核,这个配置在今天看来相当寒酸,但恰好能帮你摸到 Flink 本地模式的最低运行边界。下面我按「环境搭起来 → WordCount 跑通 → 数据流接上 → 坑填平」的顺序拆一遍,每一步都给出可复现的命令和参数说明。
2. 环境搭建与 WordCount 批处理:从零到 JAR 包提交
2.1 Flink 本地模式安装与启动参数
实验里用的是 Flink 本地模式,不涉及集群,适合开发调试。下载二进制包后解压到/usr/local/flink,然后改两个地方:conf/flink-conf.yaml里的jobmanager.rpc.address设成localhost,taskmanager.numberOfTaskSlots设成 4(跟虚拟机 4 核对齐)。启动命令是./bin/start-cluster.sh,跑完之后用jps应该能看到StandaloneSessionClusterEntrypoint和TaskManagerRunner两个进程。Web UI 默认监听 8081 端口,浏览器打开http://localhost:8081能看到 TaskManager 的可用 slot 数量。这里有个容易忽略的点:虚拟机内存只有 2048MB,Flink 默认的 JVM 堆参数可能直接把内存吃满,建议在conf/flink-conf.yaml里把jobmanager.heap.size和taskmanager.heap.size都压到512m,否则启动时容易报OutOfMemoryError。
# 解压并进入 Flink 目录 tar -xzf flink-1.13.0-bin-scala_2.11.tgz cd flink-1.13.0 # 修改配置文件,降低堆内存适配 2GB 虚拟机 sed -i 's/jobmanager.heap.size: 1024m/jobmanager.heap.size: 512m/' conf/flink-conf.yaml sed -i 's/taskmanager.heap.size: 1024m/taskmanager.heap.size: 512m/' conf/flink-conf.yaml # 启动本地集群 ./bin/start-cluster.sh # 验证进程和端口 jps | grep -E "Entrypoint|TaskManager" netstat -tlnp | grep 8081参数说明:jobmanager.heap.size控制 JobManager 的 JVM 堆上限,本地开发 512MB 足够;taskmanager.numberOfTaskSlots决定单个 TaskManager 能并行跑多少个 subtask,设成 4 是为了跟虚拟机的 4 核匹配,设太大反而会因为资源争抢导致任务排队。start-cluster.sh本质上是分别启动 JobManager 和 TaskManager 两个 Java 进程,没有守护进程化,关掉终端就没了,所以调试期间别关窗口。
2.2 Maven 项目结构与 Flink 依赖引入
IDEA 里新建 Maven 项目,pom.xml里必须引入flink-java、flink-streaming-java和flink-clients三个依赖,版本号跟 Flink 安装包保持一致。实验报告里说「Idea 里面引用 flink 报错」,九成是因为版本号写成了1.13.0但 Scala 二进制版本对不上,或者只引了flink-java没引flink-clients,导致StreamExecutionEnvironment找不到。另一个常见原因是 Maven 中央仓库下载超时,依赖只下了一半,IDEA 索引到残缺的 JAR 包就报红。解决办法是在settings.xml里配阿里云镜像,这个后面避坑章节会细说。
<dependencies> <!-- Flink Java API,版本必须与安装包一致 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-java</artifactId> <version>1.13.0</version> </dependency> <!-- 流处理 API,WordCount 和 NC 数据流都依赖它 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_2.11</artifactId> <version>1.13.0</version> </dependency> <!-- 客户端依赖,提交 JAR 包到集群时需要 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients_2.11</artifactId> <version>1.13.0</version> </dependency> </dependencies>注意flink-streaming-java和flink-clients的 artifactId 带了_2.11后缀,这是 Scala 二进制版本标识。如果你的 Flink 安装包是scala_2.12版本,这里必须改成_2.12,否则运行时会报NoSuchMethodError。很多教程只写flink-streaming-java不带后缀,在 Maven 里能解析到默认版本,但跟本地安装的 Flink 对不上,提交 JAR 包时直接抛ClassNotFoundException。
2.3 WordCount 批处理代码与打包提交
WordCount 的逻辑很直白:读文本文件,按行切词,给每个词记 1,然后按词分组求和。Flink 的flatMap负责切词并输出Tuple2<String, Integer>,groupBy(0)按第一个字段分组,sum(1)对第二个字段累加。最后用writeAsText把结果写到本地文件,或者直接print()打到标准输出。打包命令是mvn clean package,生成的 JAR 包在target目录下。提交到 Flink 的命令是./bin/flink run -c com.example.WordCount WordCount.jar,-c参数指定主类全限定名,不写的话 Flink 会从 JAR 包的 manifest 里找,但 manifest 经常因为 Maven 插件配置问题丢失主类信息,所以建议显式指定。
public class WordCount { public static void main(String[] args) throws Exception { // 批处理环境,本地模式自动识别 ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); // 读取本地文本文件,也可以换成 HDFS 路径 DataSet<String> text = env.readTextFile("/home/user/input/words.txt"); DataSet<Tuple2<String, Integer>> counts = text .flatMap(new FlatMapFunction<String, Tuple2<String, Integer>>() { @Override public void flatMap(String line, Collector<Tuple2<String, Integer>> out) { // 按空格切词,转小写避免大小写重复计数 for (String word : line.toLowerCase().split("\\s+")) { if (!word.isEmpty()) { out.collect(new Tuple2<>(word, 1)); } } } }) .groupBy(0) // 按单词分组 .sum(1); // 对计数累加 counts.writeAsText("/home/user/output/wordcount-result", FileSystem.WriteMode.OVERWRITE); env.execute("Flink Batch WordCount"); } }参数说明:readTextFile支持本地路径和 HDFS 路径,本地路径要写绝对路径,相对路径在 Flink 集群里会解析到 TaskManager 的工作目录,大概率找不到文件。FileSystem.WriteMode.OVERWRITE表示每次运行覆盖输出目录,不加这个参数第二次运行会报FileAlreadyExistsException。env.execute()是触发执行的入口,不调用的话前面所有算子都不会跑,这是 Flink 惰性求值机制决定的,新手经常写完sum就以为完事了,结果什么都没输出。
3. 数据流词频统计:NC 模拟实时数据源与流处理程序
3.1 NC 工具的使用与数据流模拟原理
nc(netcat)是 Linux 自带的网络工具,实验里用它模拟一个持续发送单词的 TCP 服务端。命令是nc -lk 9999,-l表示监听模式,-k表示保持连接不退出,9999是端口号。跑起来之后终端会阻塞,你每敲一行回车,这行文本就通过 TCP 发给所有连上来的客户端。Flink 的SocketTextStream作为客户端去连localhost:9999,就能源源不断收到数据。这里的关键是-k参数,不加的话 nc 在第一个客户端断开后就退出了,Flink 任务会报Connection reset by peer。实验报告里说「NC 程序运行后程序没有输出」,最常见的原因就是忘了加-k,或者 Flink 程序连的是127.0.0.1但 nc 监听在0.0.0.0之外的网卡上。
# 启动 nc 监听 9999 端口,-k 保持连接,-l 监听模式 nc -lk 9999 # 另开一个终端验证端口是否在监听 netstat -tlnp | grep 9999 # 手动输入测试数据,每行回车发送 hello flink hello world flink streaming注意nc的-k参数在不同 Linux 发行版上行为略有差异,Ubuntu 自带的netcat-openbsd支持-k,但有些精简版系统里的netcat-traditional不支持,需要换成ncat或者用socat替代。验证方法是nc -h看帮助信息里有没有-k选项。如果没有,可以用while true; do nc -l 9999; done循环启动,但每次客户端断开后端口会短暂释放,Flink 任务可能因为重连间隔错过数据。
3.2 流处理 WordCount 代码与窗口计算
流处理的 WordCount 跟批处理最大的区别是数据无界,不能等所有数据到齐再分组求和,必须用窗口把无限流切成有限块。实验里用的是滚动窗口(Tumbling Window),每 5 秒统计一次这 5 秒内收到的单词频率。代码结构是socketTextStream→flatMap→keyBy→window→sum。keyBy(0)按单词分组,window(TumblingProcessingTimeWindows.of(Time.seconds(5)))定义 5 秒滚动窗口,sum(1)在窗口内累加。最后print()把结果打到 TaskManager 的标准输出,在 Flink Web UI 的 TaskManager 页面能看到。
public class StreamWordCount { public static void main(String[] args) throws Exception { // 流处理环境,并行度设为 1 方便观察输出顺序 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); // 连接 nc 数据源,host 和 port 从命令行参数取 DataStream<String> stream = env.socketTextStream("localhost", 9999); DataStream<Tuple2<String, Integer>> counts = stream .flatMap(new FlatMapFunction<String, Tuple2<String, Integer>>() { @Override public void flatMap(String line, Collector<Tuple2<String, Integer>> out) { for (String word : line.toLowerCase().split("\\s+")) { if (!word.isEmpty()) { out.collect(new Tuple2<>(word, 1)); } } } }) .keyBy(0) // 按单词分组 .window(TumblingProcessingTimeWindows.of(Time.seconds(5))) // 5 秒滚动窗口 .sum(1); // 窗口内累加 counts.print().setParallelism(1); env.execute("Flink Streaming WordCount"); } }参数说明:setParallelism(1)把并行度压到 1,是为了让print()的输出按时间顺序排列,方便调试。生产环境不会这么设,但实验阶段并行度大于 1 时多个 subtask 的输出会交错,根本看不清哪个词频对应哪个窗口。TumblingProcessingTimeWindows用的是处理时间,也就是 TaskManager 的系统时钟,不是事件时间。如果 nc 发送数据的速度不均匀,窗口边界可能会切在两条数据中间,导致统计结果跟预期有偏差。想更精确的话可以换成TumblingEventTimeWindows,但需要给数据流分配时间戳和水位线,实验阶段没必要上这么复杂。
3.3 打包部署与 Web UI 监控
流处理程序打包跟批处理一样用mvn clean package,提交命令也是./bin/flink run -c com.example.StreamWordCount StreamWordCount.jar。提交之后 Flink Web UI 的 Running Jobs 页面会出现一个任务,点进去能看到算子链和每个 subtask 的状态。实验报告里说「在 flink 控制台查看输出情况」,指的就是 Web UI 里 TaskManager 的 Stdout 标签页。但注意print()的输出是写到 TaskManager 进程的标准输出,不是 JobManager 的日志,所以要去 TaskManager 页面找。如果 Web UI 里看不到输出,先确认 nc 那边有没有真的在发数据,再检查 Flink 任务的并行度和 slot 分配是否正常。
# 打包 mvn clean package # 提交到 Flink 本地集群 ./bin/flink run -c com.example.StreamWordCount target/StreamWordCount-1.0.jar # 查看运行中的任务列表 ./bin/flink list # 取消任务(如果需要) ./bin/flink cancel <job-id>flink run提交后会返回一个 JobID,用flink list能看到所有运行中的任务。如果任务提交后立刻变成 FAILED 状态,去log/flink-*-taskmanager-*.log里搜Exception,最常见的报错是java.net.ConnectException: Connection refused,说明 nc 没启动或者端口不对。另一个高频报错是NumberFormatException,通常是因为socketTextStream的端口参数写成了字符串但没做校验,或者命令行参数解析时把 host 和 port 搞反了。
4. 避坑与排查:三个真实翻车现场的血泪经验
4.1 IDEA 引用 Flink 报错:版本对齐与依赖完整性
现象:IDEA 里import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;这行直接标红,或者编译时报cannot resolve symbol。
原因:三种可能。一是pom.xml里 Flink 依赖的版本号跟本地安装的 Flink 二进制包不一致,比如本地是 1.13.0 但 pom 里写了 1.14.0,Maven 下载的 JAR 包跟运行时类库冲突。二是只引了flink-java没引flink-streaming-java,流处理相关的类找不到。三是 Maven 中央仓库下载中断,本地仓库里只有.pom文件没有.jar文件,IDEA 索引到残缺依赖就报红。
解决:先确认本地 Flink 版本,./bin/flink --version看输出。然后pom.xml里所有 Flink 依赖的版本号统一改成这个版本。如果还是报红,去~/.m2/repository/org/apache/flink/目录下检查对应的 JAR 包是否存在且大小正常,残缺的删掉重新mvn clean compile。最后在 IDEA 里右键项目 → Maven → Reload Project,强制刷新依赖索引。
4.2 Maven 打包慢:阿里云镜像配置与离线模式
现象:mvn clean package跑了十几分钟还在 Downloading,或者卡在某个依赖上不动。
原因:Maven 默认从中央仓库repo.maven.apache.org下载,国内网络访问不稳定,大文件经常断流。实验报告里说「可以在 maven 的 config 里面部署阿里云的镜像,下载速度很快」,这是最直接的解法。
解决:在~/.m2/settings.xml的<mirrors>标签里加阿里云镜像。如果settings.xml不存在就新建一个,内容如下。配好之后第一次打包还是会下载依赖,但速度能从几十 KB/s 提到几 MB/s。后续再打包时 Maven 会走本地仓库缓存,基本秒完。
<mirrors> <mirror> <id>aliyunmaven</id> <mirrorOf>*</mirrorOf> <name>阿里云公共仓库</name> <url>https://maven.aliyun.com/repository/public</url> </mirror> </mirrors>注意<mirrorOf>*</mirrorOf>表示拦截所有仓库请求,包括 Flink 自己的仓库。如果某些 Flink 快照版本在阿里云镜像里没有,会报Could not find artifact,这时候把*改成central只拦截中央仓库,Flink 快照仓库还是走官方源。
4.3 NC 程序无输出:连接建立与数据发送的时序问题
现象:nc 启动了,Flink 任务也提交了,但 Web UI 里看不到任何词频输出,TaskManager 日志也没有print结果。
原因:最常见的是 nc 启动时没加-k,Flink 的SocketTextStream连上去之后 nc 在第一次数据发送完就退出了,Flink 任务因为连接断开而卡在重连状态。另一个原因是 nc 监听在 IPv6 地址上,而 Flink 连的是 IPv4 的localhost,两者协议栈不匹配。还有一种情况是 Flink 任务的并行度大于 1,但 nc 只支持单个客户端连接,多余的 subtask 连不上就一直等待。
解决:启动 nc 时强制加-k,并用-s 0.0.0.0指定监听所有网卡的 IPv4 地址。Flink 程序里socketTextStream的 host 写127.0.0.1而不是localhost,避免 DNS 解析到 IPv6。并行度设成 1,因为SocketTextStream本身不支持并行读取,多个 subtask 连同一个 nc 端口只有第一个能收到数据。验证方法是先用telnet 127.0.0.1 9999手动连一下,能收到 nc 发的数据说明链路通了,再提交 Flink 任务。
4.4 提交 JAR 包时 ClassNotFoundException 的排查路径
现象:flink run提交后任务立刻失败,日志里抛java.lang.ClassNotFoundException: com.example.WordCount。
原因:JAR 包的 manifest 里没有主类信息,或者-c参数指定的类名跟实际包名不一致。Maven 默认的maven-jar-plugin不会自动写Main-Class到 manifest,需要手动配maven-assembly-plugin或者maven-shade-plugin打 fat JAR。
解决:最省事的办法是提交时显式加-c参数指定主类全限定名,比如./bin/flink run -c com.example.WordCount WordCount.jar。如果还是报 ClassNotFound,用jar tf WordCount.jar | grep WordCount.class确认类文件在 JAR 包里的路径跟-c参数一致。注意包名大小写敏感,com.example和com.Example是两个不同的包。
5. 进阶技巧:从实验代码到可复用 Flink 作业的四个改造点
实验代码跑通只是起点,真正要把它变成能复用的作业,还得做四件事。第一是把硬编码的参数抽出来,socketTextStream的 host 和 port 从args数组读,窗口大小从ParameterTool读,这样同一份 JAR 包能跑在不同环境。第二是给流处理程序加 Checkpoint,env.enableCheckpointing(5000)每 5 秒做一次状态快照,任务失败重启后能从上次快照恢复,不会丢窗口中间的数据。第三是把print()换成addSink,写到 Kafka 或者 MySQL,print()只适合调试,生产环境的标准输出会被 TaskManager 日志淹没。第四是给flatMap里的切词逻辑加异常捕获,遇到空行或者特殊字符时不要让整个任务挂掉。
// 从命令行参数读取配置 ParameterTool params = ParameterTool.fromArgs(args); String host = params.get("host", "127.0.0.1"); int port = params.getInt("port", 9999); int windowSize = params.getInt("window", 5); // 开启 Checkpoint,间隔 5 秒 env.enableCheckpointing(5000); env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); // 带异常保护的 flatMap .flatMap(new FlatMapFunction<String, Tuple2<String, Integer>>() { @Override public void flatMap(String line, Collector<Tuple2<String, Integer>> out) { try { for (String word : line.toLowerCase().split("\\s+")) { if (!word.isEmpty() && word.matches("[a-z]+")) { out.collect(new Tuple2<>(word, 1)); } } } catch (Exception e) { // 记录异常但不中断任务 System.err.println("Parse error: " + line); } } })ParameterTool.fromArgs(args)是 Flink 自带的参数解析工具,支持--host 127.0.0.1 --port 9999这种格式,提交时写在 JAR 包路径后面就行。enableCheckpointing(5000)的参数是毫秒,5000 表示 5 秒一次。CheckpointingMode.EXACTLY_ONCE保证每条数据只被处理一次,但会稍微增加延迟,对词频统计这种场景完全够用。word.matches("[a-z]+")过滤掉数字和标点符号,避免hello,和hello被当成两个词。
验证改造是否成功的方法:先用nc -lk 9999发一批数据,然后在 Flink Web UI 里看 Checkpoint 页面有没有成功的快照记录。再手动 kill 掉 TaskManager 进程,等几秒后 Flink 会自动重启任务并从最近的 Checkpoint 恢复,之前窗口里统计的词频不会丢。这个验证过程我每次改完流处理作业都会走一遍,因为 Checkpoint 配置写错的话,任务表面上在跑,实际上恢复时直接抛CheckpointException,等到线上出问题再查就晚了。从那以后我每次提交 Flink 作业前都强制走一遍「发数据 → 看 Checkpoint → kill TaskManager → 验证恢复」的流程,希望帮到你。
本文还有配套的精品资源,点击获取