ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Flink实战:从环境搭建到WordCount与流式词频统计

2026/9/19 4:46:31 拓冰建站 浏览量
Flink实战:从环境搭建到WordCount与流式词频统计 简介面向大数据课程实验的Flink初级编程实践报告完整覆盖WordCount程序开发与基于NC模拟数据流的实时词频统计两大任务。文档从Linux环境安装IntelliJ IDEA、Flink与Maven起步逐步展示Java代码编写、Maven打包JAR、提交Flink运行及控制台输出查看帮助初学者打通从本机调试到集群部署的完整链路。实验内容贴合入门者需求既有批处理静态统计也有流式实时计数能直观理解Flink在两种场景下的运行机制。同时针对IDEA引入Flink依赖报错、Maven打包过慢含配置阿里云镜像加速、NC命令无实时输出等典型问题给出了具体解决思路与排查方法。报告对每个环节的操作要点均进行了标注配合截图展示关键配置与运行结果适合按步骤复现也便于作为课堂实验的支撑材料。资源为1个docx文档压缩包约2.46MB内含实验环境说明、步骤截图、运行结果与问题记录便于对照复盘。已有5153人学习适合需要完成大数据实验报告或希望通过具体案例入门Flink流处理开发的学习者参考。1. Flink 编程实践的起点从一次课程实验说开去Flink 的入门门槛往往不在 API而在“从代码到运行”的链路写好的程序要能编译、能打包、能提交、能输出。很多人在 IDE 里跑通了 WordCount一到命令行提交 JAR 就卡住。这个实验把两个经典场景串在一起——批式的 WordCount 和基于 nc 的流式词频统计正好覆盖了 Flink 开发闭环里最容易出问题的几个环节。课程实验环境是 Windows 上的 VirtualBox 跑 Ubuntu这说明它不依赖特定发行版任何能装 JDK 的机器都可以照做。适合第一次接触 Flink 的学生也适合想快速排查“编包跑流程”问题的一线开发。2. 搭建 Flink 开发环境JDK、Maven 与 IDEA 的三方配合实验开始前推荐先把环境问题一次性解决。Flink 本身是用 Java 写的运行作业时要依赖 JDK而构建工具 Maven 负责拉取依赖和打包。这里不是简单“装三个软件”而是要让三者版本匹配、仓库可用否则后面每一步都可能出现玄学报错。2.1 Flink 安装与本地集群启动在 Ubuntu 上先确认 Java 版本。Flink 1.9 到 1.14 支持 JDK 8/11官方压缩包解压即用。常见做法是tar -xzf flink-1.13.2-bin-scala_2.12.tgz cd flink-1.13.2 ./bin/start-cluster.sh启动后访问http://localhost:8081能看到 JobManager 的 Web 界面就说明成功。start-cluster.sh会默认启动一个 JobManager 和若干个 TaskManager这种 stand-alone 模式适合实验不需要额外安装 Hadoop。要注意的是实验报告里写的“Flink”是纯 Flink不要为了凑 HDFS 而强行配置 Hadoop 依赖那样只会增加无意义的启动耗时。2.1.1 本地集群的内存与并行度设置实验环境本机 8GB 内存虚拟 Ubuntu 只分到 2GB。如果 Flink 默认配置的 TaskManager 内存超过 2GB启动会失败或频繁 GC。可以在conf/flink-conf.yaml中调整jobmanager.memory.process.size: 512m taskmanager.memory.process.size: 512m taskmanager.numberOfTaskSlots: 4numberOfTaskSlots决定每个 TaskManager 能并行跑几个任务虚拟机的处理器数量是 4这里设成 4 可以让默认并行度与 CPU 匹配。内存改小后作业数据量不大不会影响验证。若虚拟机启动后 Web UI 显示 TaskManager 为 0多半是内存参数配得太高调低后重启即可。2.2 Maven 项目骨架与依赖管理IntelliJ IDEA 里新建 Maven 项目后第一件事是编辑pom.xml。以 Flink 1.13.2 为例核心依赖如下properties flink.version1.13.2/flink.version scala.binary.version2.12/scala.binary.version /properties dependencies dependency groupIdorg.apache.flink/groupId artifactIdflink-java/artifactId version${flink.version}/version scopeprovided/scope /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java_${scala.binary.version}/artifactId version${flink.version}/version scopeprovided/scope /dependency /dependencies这里用了provided作用域因为 Flink 集群自带这些库打包时不需要打进 JAR。实验里如果用 IDEA 直接运行provided依赖也能在 IDE 的 classpath 中找到只是打包时不会包含。Maven 的传递依赖会拉入 flink-core、flink-runtime 等和集群环境完全一致。如果 IDEA 运行时报NoClassDefFoundError可以临时把 scope 改为compile本地验证完再改回来。不过更好的做法是在 Run Configuration 的 VM options 里加上-Dflink.lib.path/opt/flink/lib让 IDE 使用集群自带的 Flink 库从而避免 jar 冲突。2.3 用阿里云镜像解决 Maven 下载慢课程实验中最拖时间的往往不是写代码而是 Maven 下载依赖。Maven 默认中央仓库在国外容易卡。解决办法是在~/.m2/settings.xml中配置镜像mirror idaliyunmaven/id mirrorOfcentral/mirrorOf name阿里云公共仓库/name urlhttps://maven.aliyun.com/repository/central/url /mirror配置完后IDEA 需要重新导入项目可以在 Settings 里设置 Maven 的 user settings file 指向这个文件。之后依赖下载速度会有明显提升。注意mirrorOf别写成*否则会把私服也代理掉也许会引入不期望的构建行为。2.4 IDEA 引用 Flink 报错的快速定位实验报告里提到“IDEA 里面引用 flink 报错”。这个现象通常有两种原因一是 Maven 没有正确导入依赖二是 JDK 版本不匹配。先从右侧 Maven 面板点刷新按钮确认依赖列表里有没有flink-java。如果还报红检查 Project Structure 里的 Project SDK 是否选了 JDK 8 或 11。还有一点IDEA 2020 以后默认用 Maven 3 的 wrapper可能指定了旧版建议在 Settings 里换成系统 Maven。下表总结环境阶段常见的错误和对应的处理方法报错特征常见原因修复命令或操作Cannot resolve symbol ExecutionEnvironmentMaven 依赖未导入重新导入 Maven 项目java.lang.NoClassDefFoundError: org/apache/flink/api/common/...Flink 相关包被漏打包使用 maven-shade-pluginConnectException: Connection refused访问 8081JobManager 未启动执行bin/start-cluster.shMaven 下载速度慢中央仓库网络差配置阿里云镜像3. 批式 WordCount从代码到 JAR 再到 Flink 提交WordCount 是 Flink 的 Hello World但真正在生产里跑一遍比在本地 main 函数里跑要多两步打包和命令行提交。这一节用最直接的方式完成整个闭环。3.1 选择 DataSet 还是 DataStreamFlink 曾经有独立的 DataSet API批处理和 DataStream API流处理。自 Flink 1.12 开始官方明确 DataStream API 作为统一 APIDataSet 被逐步弱化。实验中的 WordCount 如果用旧教材会看到ExecutionEnvironment和DataSet但新版本里建议直接写StreamExecutionEnvironment然后通过配置启用批执行模式。两种写法都能在集群上跑区别在于流式 API 更贴近后续流式词频统计迁移成本更低。3.2 编写一个可打包的 WordCount 程序下面代码基于 DataStream API批模式运行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 WordCountBatch { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setRuntimeMode(RuntimeExecutionMode.BATCH); DataStreamString text env.fromElements( hello flink flink, hello hadoop, flink streaming); DataStreamTuple2String, Integer counts text .flatMap(new Tokenizer()) .keyBy(0) .sum(1); counts.print(); env.execute(WordCount Batch); } public static class Tokenizer implements FlatMapFunctionString, Tuple2String, Integer { Override public void flatMap(String value, CollectorTuple2String, Integer out) { String[] words value.toLowerCase().split(\\W); for (String word : words) { if (word.length() 0) { out.collect(new Tuple2(word, 1)); } } } } }参数说明env.setRuntimeMode(RuntimeExecutionMode.BATCH)让 DataStream API 以批模式执行keyBy(0)按 Tuple2 第一个字段即单词分组sum(1)对第二个字段累加。打印时print()默认是在 TaskManager 的 stdout 中输出并行度大于 1 时顺序会乱但实验阶段没有影响。3.3 打包成可提交的 JARFlink 集群内部带有 Flink 核心库但程序里如果用到自定义函数必须把这类类打进包。常见做法是使用 maven-shade-plugin 生成 fat jar。在pom.xml的buildplugins中加入plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.2.4/version executions execution phasepackage/phase goalsgoalshade/goal/goals configuration transformers transformer implementationorg.apache.maven.plugins.shade.resource.ManifestResourceTransformer mainClassWordCountBatch/mainClass /transformer /transformers /configuration /execution /executions /plugin然后执行mvn clean package -DskipTests-DskipTests跳过测试生成的文件在target/目录下形如flink-demo-1.0-SNAPSHOT.jar。注意插件里mainClass指定的是入口类这样执行flink run时可以不写-c参数但为了习惯我一般都会显式写-c。3.4 提交到 Flink 并读取结果提交命令flink run -c WordCountBatch target/flink-demo-1.0-SNAPSHOT.jar如果 Flink 部署在虚拟机需要先确认flink命令在 PATH 中或使用绝对路径./bin/flink。运行成功后控制台会输出作业执行计划日志中会看到提交成功的提示。词频结果在作业的 TaskManager 日志里位置在 Flink 安装目录的log/下可以用grep搜索。常用提交参数作用示例-c指定主类全限定名-c WordCountBatch-m指定 JobManager 地址-m localhost:8081-p设置作业并行度-p 2-d后台运行detached-d4. 流式词频统计用 nc 模拟实时数据流如果说 WordCount 是热身那么基于nc的流式词频统计就是第一次接触“无界流”。数据不再是静态集合而是持续到达的事件。这一节会演示如何用 netcat 充当数据源并写出对应的 Flink 流作业。4.1 nc 命令最简单的数据发生器Ubuntu 自带 netcat简称nc。下面的命令在 9999 端口开启一个服务端持续监听并输出收到的内容nc -lk 9999-l表示监听模式-k表示接受连接后不退出继续等待下一个客户端。默认使用 TCP 协议。启动后在另一个终端里输入hello flink并按回车这些内容就会通过网络发送到监听端。用这种方式模拟一个会持续产生文本的流源非常方便。4.2 编写实时词频统计程序Flink 程序要从 socket 上读取数据流每来一条文本就分词并更新累加器。完整代码如下import org.apache.flink.api.common.typeinfo.Types; 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 StreamWordCount { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); DataStreamString text env.socketTextStream(localhost, 9999); DataStreamTuple2String, Integer counts text .flatMap((String line, CollectorTuple2String, Integer out) - { for (String word : line.toLowerCase().split(\\W)) { if (word.length() 0) { out.collect(new Tuple2(word, 1)); } } }) .returns(Types.TUPLE(Types.STRING, Types.INT)) .keyBy(0) .sum(1); counts.print(); env.execute(Stream WordCount); } }关键点socketTextStream返回一个不会结束的 DataStream所以不能用批模式运行。setParallelism(1)是为了让print()输出按顺序显示否则多线程打印会交错。returns()告诉 Flink Lambda 表达式返回的类型这是 Java 泛型擦除带来的必要操作。4.2.1 为什么这里用keyBysum而不是groupByreduce批式 API 里可以用groupBy后接reduce但流式 API 的核心是keyBy。它把流按 key 分区每个分区内的数据会经过有状态的算子。sum(1)底层是一个增量聚合函数每来一条记录就更新一次状态不需要缓存整条数据因此能应对无限流。这也是 Flink 流处理与批处理在编程模型上的本质差异。4.3 运行与验证先启动 nc再提交作业操作顺序有讲究# 终端1启动数据源 nc -lk 9999 # 终端2提交 Flink 作业 flink run -c StreamWordCount target/flink-demo-1.0-SNAPSHOT.jar提交后回到终端1输入几行单词hello flink hello hadoop flink flink每输入一行Flink 控制台终端2的 stdout会立刻打印当前累计词频。如果终端2看不到输出需要去 Web UI 看 TaskManager 的日志。地址是http://虚拟机IP:8081进入 Task Managers → Stdout 就能看到 print 的内容。实验报告里写“在 flink 控制台查看输出”通常指的就是这个页面而不是提交作业的那个 shell。4.4 并行度与输出顺序的坑实验里如果只看到nc启动但没有输出一种情况是 Flink 作业默认并行度等于 TaskManager 的槽位数。比如虚拟机配置了4个处理器默认并行度可能就是2或4print()输出会分散到多个 TaskManager 的日志每个线程独立打印看起来像“没有输出”。这时把env.setParallelism(1)写死或者提交时加-p 1能最快看到聚合结果。生产环境中当然不能为了看日志就牺牲并行度而是要在下游接 Kafka 或文件 sink。检查点操作预期结果nc 是否监听ss -tlnp | grep 9999LISTEN状态Flink 作业状态Web UI Jobs 页面RUNNING是否有数据进入Web UI 对应 Task 的 Bytes received持续增长print 输出位置Task Managers → Stdout出现词频记录5. 让 Flink 作业在排错中跑得更稳常见问题与调试技巧实验最后把最容易卡住人的三个问题集中拆一遍再给一个从实验到项目可迁移的调试思路。5.1 高频故障的处理顺序故障现象排查路径解决方案IDEA 中 Flink 类标红检查 Maven 依赖是否导入JDK 版本重新导入用 JDK 8/11Maven 打包慢查看下载日志确认仓库 URL配置阿里云镜像并更新 settingsnc 启动后 Flink 无输出确认并行度、检查 print 位置设置setParallelism(1)去 Web UI 的 TaskManager stdout 查看ClassNotFound异常未打 fat jar配置 maven-shade-plugin5.2 用日志和 Web UI 定位数据是否真正到达当“无输出”出现时先别急着改代码。打开 Web UI 的 Running Jobs点击作业名在 Task 的 Metrics 里看numRecordsIn。如果该指标一直是 0说明数据没有进入算子问题在 nc 或 socket 连接如果numRecordsIn在增长但 stdout 没有记录大概率是并行度导致输出分散。日志方面用tail -f log/flink-taskexecutor-*.log能看到抛出的异常比如端口被占用、类型不匹配。5.3 从实验作业到生产落地的三个改造点第一不用print()做 sink改为addSink写 Kafka 或 JDBCprint()主要用于本地验证。第二增加 checkpoint在env.enableCheckpointing(5000)配合 HDFS 或 RocksDB 做持久化避免程序崩溃后状态全丢。第三使用 Flink CLI 的-d参数后台运行并用flink cancel jobId停止避免 CtrlC 直接杀掉作业而无法清理资源。flink run -d -p 1 -c StreamWordCount target/flink-demo-1.0-SNAPSHOT.jar如果 nc 端断开Flink 的 socket source 会不断重试连接。这时确认 nc 进程还活着然后用ss -tlnp | grep 9999验证端口状态。把这三个问题按顺序排查完实验里的 Flink 作业基本都能跑起来。后续换成 Kafka 或文件源只是替换 source 和 sink 的问题。本文还有配套的精品资源点击获取