简介:面向大数据初学者的《实验8 Flink初级编程实践》实验报告,完整记录了在Linux环境下使用IntelliJ IDEA进行Flink开发的全过程。实验覆盖两项核心任务:一是编写WordCount程序,经Maven打包成JAR提交至Flink集群运行;二是利用nc命令模拟数据流,编写Flink程序实时统计词频,并通过Web控制台观察输出。报告还针对实际开发中常见的Idea引用Flink报错、Maven打包过慢、nc无输出等问题给出了解决方案,可作为Flink入门实践和课程实验的参考模板。资源包内共有1个docx文档,压缩后大小约2.46MB,包含环境配置、代码编写、打包运行等关键步骤的截图与文字说明。该资源已有5153人学习下载,适合正在完成大数据实验、备战期末或初次接触Flink流处理的读者使用。
1. 实验8 Flink初级编程实践:先看懂数据流,再谈跑通
实验8 Flink初级编程实践,是绝大多数人第一次接触流式计算作业。很多同学照着模板敲一遍WordCount,看到控制台刷出几行结果就宣布“会了”,结果验收时被问到“数据从哪来、算完放哪去、并行度改成2为什么乱序”直接卡壳。这门实践真正要解决的,是把Flink编程模型的完整链路亲手打通:环境、执行计划、Source、Transformation、Sink,以及最常见的运行故障。它能帮你建立对流式计算的体感,适合正在做实验作业的学生,也适合准备Flink面试前需要动手补基础的开发。文章不会停在“跑通示例”,而是把每个环节的参数、边界和坑都拆开讲。
2. 从装环境到提交作业:Flink安装配置到部署的最小闭环
2.1 部署模式怎么选:本地模式、Standalone还是YARN
做Flink实验第一件事不是写代码,而是定部署方式。常见做法是三种:IDE里直接跑本地模式、用Docker部署Standalone集群、提交到YARN或K8s。对于“实验8”这类初级实践,我一般建议站在Standalone上做,因为本地模式掩盖了部署细节,而YARN/K8s又引入了太多外部依赖,会把实验重点带偏。
| 模式 | 适合场景 | 需要额外组件 | 实验推荐度 |
|---|---|---|---|
| IDE本地模式 | 验证API逻辑、单步调试 | 无 | 调试首选 |
| Standalone集群 | 贴近真实部署、学习资源管理 | Docker或物理机 | 推荐 |
| YARN/K8s | 生产环境、弹性资源 | Hadoop/K8s环境 | 后期再碰 |
版本选择上,Flink 1.13之后的API形态基本稳定,1.17.x和1.18.x是新用户的主力版本。做实验前先确认老师指定的版本,没有指定的话就用1.17.x系列,网上案例最多,JDK 8和JDK 11都能跑。不要一上来追最新版,很多资料里的命令在新版本里改了写法,容易踩坑。
2.2 用Docker Compose起一个两节点的Standalone集群
最省事的部署方式是用官方镜像起一个JobManager加一个TaskManager。Flink 1.16之后官方镜像直接支持jobmanager和taskmanager两个命令入口,不再像老版本那样要写一长串standalone-job.sh启动脚本。下面这个docker-compose.yml是最小可用配置:
services: jobmanager: image: flink:1.17.2 ports: - "8081:8081" - "6123:6123" command: jobmanager environment: - | FLINK_PROPERTIES= jobmanager.rpc.address: jobmanager jobmanager.memory.process.size: 1024m taskmanager: image: flink:1.17.2 depends_on: - jobmanager command: taskmanager environment: - | FLINK_PROPERTIES= jobmanager.rpc.address: jobmanager taskmanager.numberOfTaskSlots: 2 taskmanager.memory.process.size: 2048m这段配置里最容易被忽略的是jobmanager.rpc.address: jobmanager,它必须指向JobManager的容器名,否则TaskManager注册不上。端口方面,8081是Web UI,6123是RPC端口,如果本机8081被占用,改映射端口时记得把ports和容器内配置一起改。TaskManager的slot数先设2,避免实验作业并行度过高导致结果乱序时不好定位。
启动后等十几秒,打开http://localhost:8081看到Task Managers面板里出现一个节点,说明集群组件通信正常。如果TaskManager一直显示失联状态,第一步先检查两台容器的网络是否在同一网段,其次是看环境变量里的FLINK_PROPERTIES有没有被docker compose config正确解析。
2.3 第一条命令跑通官方示例:验证部署成功
集群起来后不急着写业务代码,先用自带示例确认提交链路是通的。进入JobManager容器,执行:
docker exec -it flink-jobmanager-1 ./bin/flink run \ examples/streaming/WordCount.jar \ --input /opt/flink/README.txt这里刻意用了有界输入文件而不是socket,目的是让作业跑完自动退出,方便验证“提交→调度→执行→完成”这条完整路径。执行成功后,控制台会打印出类似Job has been submitted successfully,随后能在Web UI的Completed Jobs里看到作业记录和运行时长。此时再跑一次无界流版本,观察作业状态变成Running:
docker exec -it flink-jobmanager-1 ./bin/flink run \ examples/streaming/SocketWindowWordCount.jar \ --port 9000无界作业会一直显示Running,这是流处理的正常状态。很多新手看到作业不结束就以为卡住了,实际它就是在等你往socket里发数据。验证完这两条命令,部署环节就算过关了,后面所有实验代码都提交到这个集群上跑。
3. 编程模型与词频统计:从DataStream API到你的第一个作业
3.1 执行环境与编程入口:三种写法用哪个
Flink程序的入口是执行环境,常见写法有三种:StreamExecutionEnvironment.getExecutionEnvironment()、StreamExecutionEnvironment.createLocalEnvironment()、StreamExecutionEnvironment.createRemoteEnvironment()。实验里统一用第一种,它会根据提交方式自动判断:在IDE直接运行就起本地模式,用flink run提交到集群时就连接集群的JobManager。第二种强制本地起线程模拟,第三种要手动指定JobManager地址,日常不推荐手写。
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(2);要提醒一句:setParallelism不是万能的,它只设置全局默认并行度,算子可以单独覆盖。实验里最常见的问题是全局设了2,又没理解某些算子自带并行度限制,导致输出顺序完全不可控。先把并行度设为1跑通逻辑,再改大观察行为差异,这个顺序能省去大量定位时间。
3.2 用WordCount打通“读数据-转换-输出”的完整链路
下面是一个可以直接提交到Standalone集群的WordCount程序,数据源用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; public class SocketWordCount { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); // 监听本机9000端口,作为无界数据源 DataStream<String> text = env.socketTextStream("localhost", 9000); // 切割、计数、聚合 DataStream<Tuple2<String, Integer>> counts = text .flatMap((String line, org.apache.flink.util.Collector<Tuple2<String, Integer>> out) -> { for (String word : line.split("\\s+")) { if (!word.isEmpty()) { out.collect(Tuple2.of(word, 1)); } } }) .returns(Types.TUPLE(Types.STRING, Types.INT)) .keyBy(value -> value.f0) .sum(1); counts.print(); env.execute("Socket WordCount"); } }这段代码里最关键的是.returns(Types.TUPLE(...))。Flink的lambda表达式在编译后泛型信息会被擦除,没有它,flatMap输出的Tuple2类型就识别不了,IDE里跑有时没问题,打包提交就报TypeExtractionException。这个坑几乎每个实验班都会遇到,属于典型的“本地正常、集群翻车”。
提交时用打包好的jar,注意主类全限定名要写对:
./bin/flink run -c com.example.SocketWordCount /path/to/your-job.jar先在一号终端执行nc -lk 9000监听端口,再在二号终端提交作业,回到一号终端输入hello flink,第二条终端里会打印(hello,1)、(flink,1)。到这里,一条完整的“外部数据源→API转换→控制台输出”链路就走通了。
3.3 并行度、KeyBy和“乱序输出”三者什么关系
流处理里“顺序”是伪命题。实验里经常出现这种情况:socket输入a b c,print输出却是(c,1)先出现。原因很简单,keyBy会把相同key路由到同一个子任务,不同key分到不同子任务,而print算子同样有并行度,多个并行的print线程各自往stdout写,谁先抢到输出缓冲区谁就显示在前面。
这不是代码写错了,是流计算的正常行为。想要保证输出顺序,唯一的办法是把全局并行度设成1,或者给keyBy之后的所有算子单独设置并行度1。我一般会在实验报告里写清这一点:“并行度影响吞吐也影响输出顺序,调试时优先串行,压测时再调并行度。”面试时这也是高频追问点,能说出这层关系基本就算理解了。
4. 自定义数据源与数据汇:实验里最值钱的部分
4.1 自定义Data Source:SourceFunction与运行周期
Flink自带的source类型有限,fromElements、socketTextStream、fromFile都能用于实验,但想模拟带频率的真实数据,就要自己写SourceFunction。初级实验里最常用的模板是传感器数据源:
import org.apache.flink.streaming.api.functions.source.SourceFunction; public class SensorSource implements SourceFunction<String> { private volatile boolean running = true; private int counter = 0; @Override public void run(SourceContext<String> ctx) throws Exception { while (running) { counter++; String sensorId = "sensor_" + (counter % 5); double temp = 20 + Math.sin(counter) * 5; ctx.collect(sensorId + "," + temp); Thread.sleep(1000); } } @Override public void cancel() { running = false; } }写自定义Source有两条铁律。第一,cancel()方法必须能打断run()的循环,常用手段就是这里用volatile boolean running,cancel里置false,run里的wile循环自然会退出。不要靠Thread.stop()之类的手段,Flink会标记作业失败而不是正常取消。第二,ctx.collect()是唯一合法的输出方式,不要直接在run里往外部写数据,否则Flink的算子链和checkpoint机制全部失效。
实际调试时我习惯在Source里保留休眠时间参数。Thread.sleep(1000)表示每秒一条,想测试窗口计算就改成Thread.sleep(100)模拟高频数据;想观察背压就把间隔调大到5秒,看下游算子的处理节奏。这个参数在实验报告里能写出两组对比数据,老师会很认可。
4.2 自定义Data Sink:JDBC连接器与必调参数
实验要求把结果“落库”时,最简单稳定的方案不是自定义RichSinkFunction,而是用Flink官方JDBC连接器。直接写一个SinkFunction还要自己管理连接池和重试,而JdbcSink已经把这些封装好了:
import org.apache.flink.connector.jdbc.JdbcConnectionOptions; import org.apache.flink.connector.jdbc.JdbcSink; import org.apache.flink.streaming.api.functions.sink.SinkFunction; String insertSql = "INSERT INTO word_count(word, cnt) VALUES(?, ?) ON DUPLICATE KEY UPDATE cnt = ?"; SinkFunction<Tuple2<String, Integer>> sink = JdbcSink.sink( insertSql, (ps, value) -> { ps.setString(1, value.f0); ps.setInt(2, value.f1); ps.setInt(3, value.f1); }, new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() .withUrl("jdbc:mysql://localhost:3306/flink_test") .withDriverName("com.mysql.cj.jdbc.Driver") .withUsername("root") .withPassword("123456") .withBatchSize(50) .build() );JDBC连接器有两个参数必须说明。withBatchSize(50)表示攒够50条才执行一次批量写入,这个值直接影响RDS的写入压力,实验环境设50到100都合理,生产环境要压测后定。另一个是withDriverName,MySQL 8之前用com.mysql.jdbc.Driver,MySQL 8之后必须换成com.mysql.cj.jdbc.Driver,写错的话作业一直报找不到驱动类或SSL连接异常。
很多人问为什么不用自定义SinkFunction连接MySQL,我的观点是:实验目的是理解Flink编程模型,不是重复造连接池轮子。JDBC连接器足够完成“把结果写到外部系统”这个教学点,而且踩坑少。自定义SinkFunction留到进阶实验里和Redis、Elasticsearch一起写更有意义。
4.3 实验报告常问的背压是什么:一句话答清楚
背压是流处理系统里下游处理速度跟不上上游生产速度时,系统自动让上游放慢的一种机制。实验里用自定义Source模拟每秒1000条数据,下游用Thread.sleep(50)模拟慢处理,Web UI的BackPressure面板会从OK变为HIGH。这说明Flink在自动协调上下游速率,不是丢数据。回答面试时补充一句“背压是流处理设计的一部分,不是异常”,比背一段定义强得多。
5. Flink初级编程避坑指南:5条高频翻车记录
5.1 打印不出结果:并行度大于1把stdout拆碎了
现象:作业在Web UI上是Running状态,但IDE控制台或TaskManager日志里看不到print输出。原因:print()算子的默认并行度继承作业全局并行度,比如全局设了4,每个并行子任务独立打印,日志分散在4个TaskManager的stdout里,本地看不到或只看到一部分。解决:调试期把env.setParallelism(1)或单独给print设置setParallelism(1),让所有结果汇总到一个线程输出。生产环境不要依赖print看结果,应该接入日志系统或落到外部存储。
5.2 JDBC连接器异常:驱动类在提交端不在容器里
现象:作业提交几秒后失败,报ClassNotFoundException: com.mysql.cj.jdbc.Driver,或者直接报无法加载驱动。原因:依赖的scope写成了provided,flink-connector-jdbc和MySQL驱动没有打进jar,集群上自然找不到。解决:pom.xml里把JDBC连接器和驱动的scope改为默认的compile,并确认使用maven-shade插件打包。我在实验环境见过最离谱的情况是驱动包在IDE的lib目录里有,同学以为提交时也会带上,结果集群上必然报错。打包后检查jar里的BOOT-INF/lib或根目录下是否有驱动class文件,比反复提交省时间。
5.3 Sink到Hive表数据不入表:分区目录与提交时机
现象:作业执行成功,Web UI显示无异常,但Hive表里查不到新增数据。原因:Hive Sink在流式写入时会先把数据写到分区目录,等checkpoint完成才提交文件;如果作业没开启checkpoint,或者用户按批处理思维等作业结束才看表,常见结果是分区目录里有part-临时文件但表分区元数据没有更新。解决:在env上开启周期性checkpoint,例如env.enableCheckpointing(10000),然后等一个完整checkpoint周期再查Hive。另外确认写入模式不是严格一次而是至少一次,实验场景两者都能接受。这个坑在真实项目里最容易让人暴躁,因为“看起来成功了但数据就是没进表”。
5.4 本地模拟器跑通、集群上ClassNotFound
现象:IDE里运行一切正常,打包用flink run提交到Standalone就报NoClassDefFoundError或各种方法找不到。原因:Maven依赖里写了<scope>provided</scope>,这个scope的意思是“运行环境提供这个依赖”,本地IDE自动把依赖加进classpath所以能跑,提交时Flink集群没有这些类就炸了。解决:将flink-streaming-java保持provided没问题,因为集群确实自带Flink核心依赖;但第三方库比如JSON库、Kafka客户端,必须去掉provided。用maven-shade生成fat jar时,排查mvn dependency:tree里有没有遗漏。这是一条最经典的生产级坑,实验阶段踩一次比面试背十遍都管用。
针对这个坑,给出一个可直接套用的shade插件配置:
<plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.5.1</version> <executions> <execution> <phase>package</phase> <goals><goal>shade</goal></goals> <configuration> <transformers> <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"> <mainClass>com.example.SocketWordCount</mainClass> </transformer> </transformers> </configuration> </execution> </executions> </plugin>5.5 Watermark不触发窗口计算
现象:定义了事件时间和窗口,数据也源源不断进来,但窗口就是不输出结果。原因:窗口关闭依赖水位线,而水位线推进需要遇到足够新的事件事件戳;如果并行度大于1,每个并行子任务各自维护水位线,Flink按最小的那个推进全局水位线,只要某个子任务没有新数据,全局水位线就停滞,窗口永远不关。解决:初级实验先用env.setParallelism(1)排除这个因素,确认为多并行度问题后,再考虑在source上统一分配水位线,或者观察子任务水位线差值。这个点把作业从“能跑”推到“能解释”,实验报告里写清排查过程价值很高。
6. 从“跑通”到“讲得清”:验证作业正确性的三个方法
作业跑通不等于结果对。我的习惯是至少做三件事验证。第一,Web UI的指标面板对比各算子之间的Records Sent和Records Received数量,两个数字长期不一致说明算子逻辑里丢了数据。第二,把同样一份有界数据同时跑Flink和简单的批处理SQL,两者结果对不上一定有一方理解错了业务需求。第三,单并行度跑一遍逻辑,再逐步调大并行度,对比结果是否一致;不一致就检查keyBy逻辑和状态使用是否正确。
进阶一点,可以打开火焰图观察算子热点。Flink Web UI的Profiler功能能抓取JobManager和TaskManager的CPU采样,火焰图里如果某个map算子占比异常高,优先看是不是序列化开销,其次看业务逻辑里有没有无谓的字符串拼接。我见过最典型的案例是自定义Sink里用System.out.println打日志,火焰图里输出操作占了30%的CPU,换成logback之后性能立刻回升。
最后说一个让我印象最深的教训:有一次做实验,我用keyBy(sensorId)聚合数据,但sensorId上游拼错了大小写,同一传感器在结果表里分成了两条记录。检查代码逻辑查不出错,最后用SQL对了一下聚合结果才发现。从此我每次写完作业都会顺手跑一遍对照查询,这个习惯一直保留到现在。流处理程序不看结果只看运行状态,很容易被“作业Running”骗过去。希望帮到你。
本文还有配套的精品资源,点击获取