
简介本资源为大数据课程实验8「Flink初级编程实践」的完整实验报告面向正在学习大数据技术原理与应用的高校学生及Flink初学者帮助解决从环境搭建到作业提交的全流程实践问题。报告围绕两个核心任务展开一是使用IntelliJ IDEA开发WordCount程序涵盖Flink与Maven安装、Java代码编写、JAR打包及集群运行二是借助Linux自带NC程序模拟实时数据流编写Flink程序完成词频统计并部署运行。文中还记录了Idea引用Flink报错、Maven打包缓慢、NC程序无输出等典型问题的排查与解决思路并附有Flink Web控制台查看输出的方法。资源包为1个docx文档约2.46MB结构完整、步骤清晰适合对照复现实验与查漏补缺。目前已有5153人学习下载可作为大数据实验课提交与Flink入门练习的参考材料。1. Flink初级编程实践从本地环境到第一个能跑通的DataStream作业很多同学第一次接触 Flink 是在课程实验里标题写着“初级编程实践”打开一看却要装集群、配 YARN、连 Kafka直接卡在环境这一步。其实 Flink 初级编程的核心只有一件事把 Source、Transformation、Sink 这条链路在本地跑通理解算子之间的数据流转。你不需要先有一整套大数据平台一台装了 JDK 的笔记本就够。这篇笔记面向正在做 Flink 入门实验、想自己动手写第一个作业的人也适合已经会写 SQL 但没碰过 DataStream API 的后端同学。我会按“环境怎么搭、代码怎么写、参数怎么调、坑在哪”的顺序讲每一步都能直接抄。Flink 编程实践最怕的不是逻辑复杂而是环境玄学和依赖冲突先把最小可运行版本跑起来后面加 Kafka、JDBC、ClickHouse 才有意义。2. 本地跑通 Flink 作业环境、依赖与最小骨架2.1 为什么初级实践优先用本地执行模式Flink 有三种执行环境本地LocalExecutionEnvironment、远程RemoteEnvironment、以及提交到集群的 Standalone/Session 模式。初级编程实践阶段我强烈建议先用本地模式。原因很直接本地模式不需要启动 JobManager 和 TaskManager 进程代码里main方法一跑Flink 会在当前 JVM 里起一个 MiniCluster算子并行度默认等于 CPU 核数。这样你能把注意力放在算子逻辑上而不是“为什么 TaskManager 注册不上”。本地模式还有一个好处是调试方便。你可以在map、filter里直接打断点变量值看得一清二楚。一旦切到集群模式日志分散在多个节点初级阶段的排错成本会陡增。常见做法是本地模式验证逻辑再改成StreamExecutionEnvironment.getExecutionEnvironment()提交到集群。依赖方面Maven 里至少需要flink-streaming-java和flink-clients。注意 Flink 1.15 之后flink-clients被拆出来不引会报No ExecutorFactory found。版本号建议和你要部署的集群保持一致避免序列化器不兼容。2.2 用 Maven 搭一个能跑的最小工程先建一个普通 Maven 项目pom.xml里加以下依赖。这里以 Flink 1.17 为例你按自己集群版本替换即可。properties flink.version1.17.1/flink.version java.version11/java.version /properties dependencies !-- Flink 核心流处理 API -- dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version${flink.version}/version /dependency !-- 本地执行必需1.15 之后独立出来 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-clients/artifactId version${flink.version}/version /dependency /dependencies逻辑说明flink-streaming-java提供DataStream、MapFunction等核心类flink-clients提供本地执行器工厂。参数说明flink.version必须和运行环境一致混用 1.14 和 1.17 的包会出现ClassNotFoundException。Java 版本建议 11Flink 1.17 对 Java 17 支持还不完整用 17 可能遇到模块访问警告。2.3 第一个 DataStream 作业从集合到控制台下面这段代码是 Flink 初级编程实践里最经典的 WordCount 简化版数据源用集合不依赖外部系统。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 LocalWordCount { public static void main(String[] args) throws Exception { // 本地执行环境并行度默认取 CPU 核数 StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); // 初级调试建议设为 1输出顺序稳定 DataStreamString lines env.fromElements( flink programming practice, flink datastream api, flink programming ); DataStreamTuple2String, Integer counts lines .flatMap(new FlatMapFunctionString, Tuple2String, Integer() { Override public void flatMap(String value, CollectorTuple2String, Integer out) { for (String word : value.split( )) { out.collect(Tuple2.of(word, 1)); } } }) .keyBy(t - t.f0) // 按单词分组 .sum(1); // 对第二个字段累加 counts.print(); // 输出到控制台 env.execute(Local WordCount); } }逻辑说明fromElements创建一个有界流适合实验flatMap把每行拆成单词并输出(word,1)keyBy按单词分区保证相同单词落到同一算子实例sum(1)对索引为 1 的字段累加。参数说明setParallelism(1)让所有数据进同一个 subtask打印结果不会交错如果设为默认并行度print()的输出顺序会乱初学者容易误以为逻辑错了。运行后控制台会输出类似(flink,3)、(programming,2)的结果。到这一步你的 Flink 初级编程环境就算通了。3. 把 Source、Transformation、Sink 拆开练参数与算子选择3.1 Source 的三种常见写法与适用场景初级实践里 Source 通常有三种集合、文件、Socket。集合用fromElements或fromCollection适合单元测试文件用readTextFile适合批处理实验Socket 用socketTextStream适合模拟实时流。// 文件 Source路径可以是本地路径或 HDFS 路径 DataStreamString fileStream env.readTextFile(data/input.txt); // Socket Source先启动 nc -lk 9999 DataStreamString socketStream env.socketTextStream(localhost, 9999);逻辑说明readTextFile默认按行读取支持本地和分布式文件系统socketTextStream会持续监听端口每来一行就触发一次处理。参数说明文件路径在本地模式下用相对路径即可提交集群时要改成绝对路径或 HDFS 路径Socket 的 host 和 port 要和nc命令一致否则会一直重试连接。注意readTextFile在 Flink 1.17 里属于旧版 Source API新项目建议用FileSource但初级实验用旧 API 更省事不用额外引flink-connector-files。3.2 Transformation 里最该先掌握的四个算子初级编程实践不需要把算子全背下来先掌握map、filter、keyBy、reduce这四个就能覆盖大部分实验题。DataStreamString filtered socketStream .filter(line - line ! null !line.trim().isEmpty()) // 过滤空行 .map(String::toLowerCase); // 转小写 DataStreamTuple2String, Integer aggregated filtered .map(word - Tuple2.of(word, 1)) .returns(Types.TUPLE(Types.STRING, Types.INT)) .keyBy(t - t.f0) .reduce((a, b) - Tuple2.of(a.f0, a.f1 b.f1)); // 增量聚合逻辑说明filter去掉空行避免后续split报空指针map做归一化keyBy按单词分组reduce对每组做增量聚合比sum更灵活。参数说明returns显式声明元组类型因为 Lambda 的类型擦除会让 Flink 无法推断不加会抛InvalidTypesException。这是初级实践里最常见的翻车点之一。3.3 Sink 输出到控制台、文件与 JDBCprint()是最简单的 Sink但实验报告通常要求输出到文件或数据库。文件 Sink 用writeAsTextJDBC Sink 需要引flink-connector-jdbc。// 输出到文件并行度设为 1 时只生成一个文件 counts.writeAsText(output/result, FileSystem.WriteMode.OVERWRITE); // JDBC Sink 示例需先建表 counts.addSink(JdbcSink.sink( INSERT INTO word_count (word, cnt) VALUES (?, ?) ON DUPLICATE KEY UPDATE cnt ?, (ps, t) - { ps.setString(1, t.f0); ps.setInt(2, t.f1); ps.setInt(3, t.f1); }, JdbcExecutionOptions.builder() .withBatchSize(100) .withBatchIntervalMs(200) .withMaxRetries(3) .build(), new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() .withUrl(jdbc:mysql://localhost:3306/test) .withDriverName(com.mysql.cj.jdbc.Driver) .withUsername(root) .withPassword(123456) .build() ));逻辑说明writeAsText把结果写到目录并行度大于 1 时会生成多个 part 文件JdbcSink.sink用批处理方式写入减少数据库压力。参数说明withBatchSize(100)表示攒够 100 条才提交一次withBatchIntervalMs(200)表示最多等 200 毫秒withMaxRetries(3)是失败重试次数。这三个参数直接决定写入吞吐和延迟初级实验里设小一点方便观察。4. 避坑与排查Flink 初级实践里最容易翻车的五件事4.1 现象本地跑报 No ExecutorFactory found原因Flink 1.15 之后flink-clients从flink-streaming-java里拆出去了只引流处理包找不到本地执行器。解决在pom.xml里显式加flink-clients依赖版本和流处理包保持一致。4.2 现象Lambda 表达式报 InvalidTypesException原因Java Lambda 的类型擦除导致 Flink 无法推断Tuple2的泛型。解决在map或flatMap后加.returns(Types.TUPLE(Types.STRING, Types.INT))或者改用匿名内部类实现MapFunction。4.3 现象print 输出顺序乱、结果对不上原因默认并行度等于 CPU 核数多个 subtask 并行输出控制台交错。解决调试阶段env.setParallelism(1)确认逻辑正确后再调大并行度。注意keyBy之后的sum结果本身是对的只是打印顺序乱。4.4 现象Socket Source 一直连不上原因nc -lk 9999没启动或者 host 写成了容器 IP。解决先在终端执行nc -lk 9999再运行 Flink 程序如果 Flink 跑在 Docker 里host 要用宿主机 IP不能用 localhost。4.5 现象JDBC Sink 报驱动找不到原因flink-connector-jdbc不带 MySQL 驱动需要单独引mysql-connector-java。解决加 MySQL 驱动依赖并确认withDriverName写的是com.mysql.cj.jdbc.DriverMySQL 8 之后的新类名旧版com.mysql.jdbc.Driver会警告。5. 从初级作业到可复用模板并行度、水位线与检查点怎么加初级实践跑通之后下一步是让作业具备“可复用”的骨架。我一般会在模板里固定三件事并行度、水位线、检查点。并行度通过env.setParallelism或算子级setParallelism控制Source 和 Sink 的并行度可以单独设比如 Source 用 1 保证顺序中间算子用 4 提高吞吐。水位线用WatermarkStrategy配置事件时间场景下必须设否则窗口不触发。env.enableCheckpointing(5000); // 每 5 秒做一次检查点 env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(2000); env.getCheckpointConfig().setCheckpointTimeout(60000);逻辑说明enableCheckpointing(5000)开启检查点间隔 5 秒EXACTLY_ONCE保证精确一次setMinPauseBetweenCheckpoints防止检查点太密集拖慢处理setCheckpointTimeout超时则丢弃本次检查点。参数说明初级实验可以把间隔设大一点比如 10000 毫秒减少日志干扰。验证方法很简单在map里加一个计数器观察检查点触发时计数器是否回滚。如果回滚说明状态后端生效了。我自己的习惯是每写一个新算子先用并行度 1 加print验证再逐步调大并行度、加检查点。这样出问题时能快速定位是逻辑错还是配置错。希望帮到你。本文还有配套的精品资源点击获取