
简介《实验8 Flink初级编程实践》是一份面向大数据课程学习者的完整实验报告聚焦Flink流处理框架的入门操作。报告基于Windows 10主机与Linux Ubuntu虚拟机环境详实记录了IntelliJ IDEA开发WordCount程序、Maven打包JAR并提交Flink运行的全过程并利用NC工具模拟实时数据流完成词频统计的实时处理与结果查看两个任务分别覆盖批处理与流处理的基本编程范式包含环境搭建、依赖配置、打包部署和Web控制台监控等关键环节。文中还整理了IDE引用Flink报错、Maven打包缓慢、NC无输出等常见问题的排查与解决方案对正在学习大数据技术或准备相关实验的读者具有直接参考价值。资源为单个docx文档大小2.46MB内容包含实验环境、操作截图、代码与运行结果条理清晰。该报告已有5153人学习浏览适合作为Flink初级实践的对照模板与排错指南。1. 实验8 Flink初级编程实践不是“会调API”是能跑通一整个pipeline很多实训课把Flink的第一次动手放在第8个实验前面全是环境和概念的铺垫。真实场景里这个实验的翻车率比想象中高在IDE里点右键跑通WordCount离“会写Flink”还很远同一个jar扔到集群上就报ClassNotFoundException、资源不足、结果打不到预期位置。这门实验真正要解决的是三件事把一个Flink程序从工程骨架、运行环境到提交执行的链路打通把Source、Transformation、Sink三个角色在代码里的分工看清楚把“本地能跑”和“集群能跑”的差异彻底摸透。适合刚学完Java和基础的Hadoop、准备正式踩Flink这门手艺的初学者也适合一直只调API、从没看过集群作业的开发者回炉一遍。2. 安装配置到部署本地模式下把第一个WordCount跑起来2.1 环境先行本地模式、Docker还是集群网上能搜到大量“Flink安装配置到部署”的教程但很多菜鸟教程一上来就教你怎么用IDE写API忽略了环境差异。Flink的部署形态至少分三种本地单机模式Local Cluster、Docker容器部署、独立Standalone集群。实验8阶段我强烈建议先用本地模式不是Docker不好而是本地模式日志最直白、Web UI最完整出了问题看日志文件比绕一层容器少一道排查。等你在本地把WordCount跑通再谈Docker和集群不迟。Flink官方压缩包解压后bin目录下的start-cluster.sh就是本地模式的启动入口。它会拉起一个JobManager进程和至少一个TaskManager进程Web UI默认监听8081端口。启动后先别急着写代码把环境确认这一步做扎实# 解压并启动本地集群 tar -zxf flink-*.tgz cd flink-* bin/start-cluster.sh # 确认进程起来了 jps | grep -E JobManager|TaskManager # 确认Web UI能访问 curl -s http://localhost:8081 | head -5说明一下这段命令做了什么事start-cluster.sh会读取conf/flink-conf.yaml中的内存和并行度配置然后分别启动JobManager和TaskManager。jps是JDK自带的进程查看工具如果列出的JVM进程中没有TaskManager说明资源分配或JDK版本可能有问题别急着往下走。curl那一步是确认Web UI真的在服务而不是端口被占用返回了别的页面。默认并行度在conf/flink-conf.yaml里的taskmanager.numberOfTaskSlots字段控制实验8阶段保持默认1就行。2.2 用Maven快速生成工程骨架pom.xml里的版本玄学Flink的工程结构和普通Java项目不完全一样依赖需要区分作用域主类要能被shade插件正确打包flink-version和scala.binary.version必须对得上。自己手动建目录容易漏常见做法是用官方archetype生成骨架mvn archetype:generate \ -DarchetypeGroupIdorg.apache.flink \ -DarchetypeArtifactIdflink-quickstart-java \ -DarchetypeVersion1.17.1 \ -DgroupIdcom.experiment \ -DartifactIdflink-lab \ -Dversion1.0-SNAPSHOT \ -Dpackagecom.experimentarchetypeVersion就是你本地Flink的版本不要随便写一个和你集群不一致的版本号。生成后pom.xml是这个实验的关键文件我一般会精简成下面这样properties flink.version1.17.1/flink.version scala.binary.version2.12/scala.binary.version maven.compiler.source1.8/maven.compiler.source maven.compiler.target1.8/maven.compiler.target /properties dependencies dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version${flink.version}/version scopeprovided/scope /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-clients/artifactId version${flink.version}/version scopeprovided/scope /dependency /dependencies两个陷阱说明flink-streaming-java和flink-clients都必须用provided作用域意思是编译时需要、运行时不需要打进jar因为Flink的lib目录里已经有这些包。如果不用provided打出来的jar会带上大量Flink类提交到集群时轻则冲突、重则直接ClassCastException。另一个陷阱是scala.binary.versionFlink的Java工件命名里带有Scala版本后缀2.12是主流但如果你的Flink是预编译的Scala 2.11版本这里就必须改成2.11否则依赖根本拉不下来。2.3 词频统计初体验第一个Flink程序的完整代码实验8最标准的练手程序就是词频统计。它覆盖了Source、Transformation、Sink三个角色代码短但结构完整。我用fromElements作为输入源这样不依赖外部服务跑起来最干净import org.apache.flink.api.common.typeinfo.TypeHint; import org.apache.flink.api.common.typeinfo.TypeInformation; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.streaming.api.datastream.DataStreamSource; import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.util.Collector; public class WordCount { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(2); DataStreamSourceString lines env.fromElements( hello flink, flink word count, hello world); SingleOutputStreamOperatorTuple2String, Long words lines .flatMap((String line, CollectorTuple2String, Long out) - { for (String w : line.split( )) { if (!w.isEmpty()) { out.collect(Tuple2.of(w, 1L)); } } }) .returns(TypeInformation.of(new TypeHintTuple2String, Long() {})); words.keyBy(t - t.f0) .sum(1) .print(); env.execute(word-count-lab); } }这段代码是Flink初级编程的标准套路。env.fromElements把三行字符串变成数据流数据流的每个元素就是一行文本flatMap把一行拆成多个单词每拆出一个就collect一个(单词, 1)的二元组returns那行必须写因为lambda表达式在Java里会有泛型类型擦除不显式指定类型信息的话Flink在编译期可能直接抛“could not determine type”异常keyBy按单词分区sum(1)把二元组第二个字段累加print把结果打到标准输出env.execute是真正触发作业执行的地方后面详说。注意setParallelism(2)这行它把整条数据流的并行度设成2这样print也会用两个并行子任务结果顺序看起来就是乱的这是下一章的伏笔。2.4 提交与找输出flink run和IDE Run的差别在IDE里直接右键run用的是当前进程作为JobManager日志全打在IDE控制台这会造成一种假象Flink就是这么简单。实际上集群环境里提交作业的命令是下面这样mvn clean package -DskipTests /opt/flink/bin/flink run \ -c com.experiment.WordCount \ target/flink-lab-1.0-SNAPSHOT.jar-c指定主类全限定名后面跟jar包路径。提交成功后print的结果不会出现在终端而是写到TaskManager的日志文件里常见位置是Flink目录下的log文件夹# 实时看taskexecutor的标准输出 tail -100f log/*taskexecutor*.out这是我反复提醒初学者的第一个坑找不到输出不是作业挂了而是输出在一个你想不到的位置。真正判断作业跑没跑成去Web UI上看那个界面会告诉你每个算子的记录数、吞吐和背压比猜强一百倍。3. 词频统计背后的DataStream模型Source、Transformation、Sink如何衔接3.1 三角色模型和懒执行机制Flink程序为什么必须“点火”把上面的WordCount跑通以后下一件事不是学更多算子而是把Flink的编程模型看透。Flink的DataStream API可以简化为三个角色的流水线Source负责产生数据Transformation负责变换数据Sink负责输出数据。词频统计里fromElements是SourceflatMap和keyBy/sum是Transformationprint是Sink。这个模型面试必考但更重要的是理解背后的执行机制所有算子只是被“组装”成一个逻辑DAG直到你调env.execute()整个作业才会被编译、提交并真正开始跑。这就是所谓的懒执行。实操中很多人把env.execute()漏掉或者写在某个if分支里结果作业提交了却立刻结束、什么都不输出。排查思路很简单看控制台有没有打印出“Job has been submitted”这类日志没有的话检查execute()是不是真的被调到了。懒执行还意味着你如果在execute之后才对数据流做处理那部分代码永远不生效。3.2 并行子任务与算子链为什么print出来的结果是乱序的上一章留了一个伏笔setParallelism(2)之后print输出是乱的。这不是Bug。Flink的每条数据流在逻辑上是一个算子链但物理执行时会按并行度拆成多个并行子任务。setParallelism(2)把Source、flatMap、keyBy、print全部拆成两个并行实例每个实例独立运行、独立输出两个实例往同一个stdout里写顺序自然交错。确认这个机制的办法很简单在Web UI里点击正在运行的作业你会看到一张Job Graph每个节点下方标注着并行度数字例如“Source: Collection Source (2/2)”。算子链是另一个新手容易困惑的点。你会发现Job Graph里某个节点叫“Source: Collection Source - FlatMap”说明这两个算子被串在了同一个Task里中间不落网络、不落内存。这是Flink的优化满足链条件时自动合并减少序列化和网络传输。但副作用是Job Graph的粒度比代码里的算子数量少别以为“我写了六个算子就一定看到六个节点”。看执行计划时认准并行度和记录数比纠结有几个节点重要得多。3.3 提交参数与默认行为-p、-d、-m怎么配合从“IDE里跑”切换到“命令行提交”需要建立一套参数心智模型。下面这条命令把前面WordCount的提交完整化/opt/flink/bin/flink run \ -m localhost:8081 \ -c com.experiment.WordCount \ -p 2 \ -d \ target/flink-lab-1.0-SNAPSHOT.jar参数说明-m指定JobManager的地址和端口默认就是localhost:8081本地模式可以省略-p覆盖作业并行度优先级高于代码里的setParallelism也高于conf里的默认值命令行指定的值最霸道-d表示detached提交后立刻返回shell不挂在终端跟前台等待作业结束。实际工作中-p这个参数是最容易拍脑袋的并行度不是越大越好它会直接对应slot占用本地模式默认只有1个TaskManager、1个slot-p设置成4作业会一直等资源。还有个默认行为必须知道Flink默认不开启Checkpoint。意思是作业一旦异常重启状态从头计算。词频统计这种无状态作业无所谓但如果你在后面实验中加了状态或窗口就要主动开启Checkpoint否则数据丢失时你甚至没有后悔药。开启的方式是env.enableCheckpointing(5000)每5秒做一次快照实验8可以先不碰但要记住这个默认行为。最后提醒一个老教程的坑网上大量词频统计教程用的是env.socketTextStream(localhost, 9999)从Socket读数据新版Flink已经不推荐这种方式测试环境可以用DataStreamUtils或本案例的fromElements替代。照抄旧代码编译报错先查API是否被移除再去改代码别在环境上折腾半天。4. 自定义Data Source与Data Sink从模拟日志到JDBC写入的完整链路4.1 自定义Data Source一用RichSourceFunction造一个可控的模拟数据源实验8的进阶要求通常会落在自定义Source和Sink上。为什么要自定义Source因为fromElements只能测语法真实作业的数据来自Kafka、MySQL、文件或埋点SDK这些都需要实现Source接口来接入。这里我做的是一个每分钟循环生成随机日志的自定义Source它是后续接真实数据源的最小模板。import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.functions.source.RichSourceFunction; import org.apache.flink.streaming.api.functions.source.SourceContext; import java.util.Random; public class RandomLogSource extends RichSourceFunctionString { private volatile boolean running true; private final String[] keywords {flink, java, hadoop, spark}; Override public void open(Configuration parameters) throws Exception { // 可以在这里初始化数据库连接或文件句柄 } Override public void run(SourceContextString ctx) throws Exception { Random random new Random(); while (running) { String line keywords[random.nextInt(keywords.length)] _ System.currentTimeMillis(); ctx.collect(line); Thread.sleep(500); } } Override public void cancel() { running false; } }代码逻辑分三块open在算子初始化时执行一次适合放连接资源run在后台循环里不停生产数据每500毫秒发射一条cancel用于任务取消时把running置false让run中的循环退出。三个细节是血泪经验第一running必须用volatile修饰否则多线程环境下cancel的修改对run线程不可见第二Thread.sleep必须在collect之后否则背压还没形成就把自己打满第三cancel里只置标志位就够了不要关闭连接资源资源的释放放到close方法里。这个Source接入主程序就一行env.addSource(new RandomLogSource()).print()。4.2 自定义Data SinkRichSinkFunction加JDBC把结果真正落库Source是数据入口Sink是数据出口。print只把数据打到标准输出真正生产里要写MySQL、ES或HDFS。我常见的做法是写一个RichSinkFunction的子类在open里建立连接在invoke里逐条写入在close里释放连接import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.functions.sink.RichSinkFunction; import java.sql.Connection; import java.sql.DriverManager; import java.sql.PreparedStatement; public class JdbcWordSink extends RichSinkFunctionorg.apache.flink.api.java.tuple.Tuple2String, Long { private Connection conn; private PreparedStatement insertStmt; Override public void open(Configuration parameters) throws Exception { Class.forName(com.mysql.cj.jdbc.Driver); conn DriverManager.getConnection( jdbc:mysql://localhost:3306/test?useSSLfalseserverTimezoneAsia/Shanghai, root, your_password); insertStmt conn.prepareStatement( INSERT INTO word_count (word, cnt) VALUES (?, ?)); } Override public void invoke(Tuple2String, Long value, Context context) throws Exception { insertStmt.setString(1, value.f0); insertStmt.setLong(2, value.f1); insertStmt.executeUpdate(); } Override public void close() throws Exception { if (insertStmt ! null) insertStmt.close(); if (conn ! null) conn.close(); } }这个类的时序是open在subtask启动时调用一次invoke对每条数据调用一次close在subtask关闭时调用一次。invoke里不要做重活executeUpdate如果不开批量模式单条写入会非常慢生产环境要改成executeBatch或使用官方JdbcSink的批量接口。URL里的serverTimezoneAsia/Shanghai是新手最容易漏的不加的话MySQL 8的驱动会报时区错误看起来像连接失败实际是参数问题。4.3 JDBC连接器异常、Hive Sink不入表和CDC Pipeline三个不该混进来的边界这里把几个热搜里的高频问题一并讲清楚。先讲JDBC连接器异常最常见的是“No suitable driver found”现象是open阶段直接抛SQLException原因通常是Class.forName里写的驱动类名和实际依赖不匹配MySQL 5用com.mysql.jdbc.DriverMySQL 8用com.mysql.cj.jdbc.Driver。第二个常见现象是“Communications link failure”可能是数据库地址不可达也可能是URL里没带时区参数。排查顺序固定先确认驱动类能被加载再确认数据库从TaskManager所在机器能连通最后看URL参数完整性。再说“Flink Sink Hive表数据不入表”。这个现象经常出现在有人把Hive当普通JDBC写、或者在Flink SQL里建Hive表但开了流式模式却没配分区。实验8阶段不要碰Hive Sink它依赖独立的HiveCatalog和Hive方言解决的问题和Java DataStream API是两条线。类似地还有“Flink CDC Pipeline部署”CDC是把数据库变更流接入Kafka等消息系统它属于Flink的另一套生态你连JobGraph都还没吃透就去追CDC只会被一堆deployment日志绕晕。数据血缘例如OpenMetadata抓取血缘更是数据治理层面的事这些方向都没有错但不是“实验8 Flink初级编程实践”要回答的问题。提示判断一个需求是否属于初级实践就看它是否依赖额外Catalog或外部同步框架。只用DataStream API和连接器能写完的就是这一阶段的主题需要引入Kafka Connect、Hive方言或CDC框架的尽管名字里带Flink也不要在这个实验里混学。5. Flink初学避坑指南5个常见的翻车现场与排查路径5.1 打包上集群就ClassNotFoundshade插件和依赖作用域现象IDE里run一切正常mvn package之后flink run提交几秒钟后报java.lang.ClassNotFoundException或NoClassDefFoundError有时候提到的类是你自己写的Main类有时候是第三方库的类。原因多数是maven-shade-plugin没配或者配了但flink依赖没有用provided作用域导致jar冲突。如果flink-streaming-java被以compile作用域打进了jar提交时就会和集群lib目录里的Flink类发生版本冲突报出来的异常五花八门。解决pom.xml里加上maven-shade-plugin并把所有flink前缀的依赖都设为provided第三方连接器依赖保持默认compile。打出来的jar最好控制在几MB内超过50MB就要怀疑是不是把Flink打进去了。5.2 词频统计结果“乱序”print不是有序输出这是误解不是Bug现象print出来的结果一会儿输出hello一会儿输出world而且统计值看起来是乱跳的很多人以为是flatMap写错了。原因并行度大于1时print算子有多个并行子任务每个子任务独立往stdout里写数据操作系统层面对多个线程的输出没有顺序保证。这是Flink分布式执行的正常表现不是代码问题。解决确认数据本身正确用Web UI看Print Sink的“Records Received”数字在涨就说明数据在流通。想看有序输出把print单独设成setParallelism(1)或者把整体并行度设为1来测试逻辑。但要注意这只适合看结果生产环境的并行度该调高还得调高。5.3 JDBC连接器异常No suitable driver、时区与连接超时现象作业提交后open阶段抛异常有的是No suitable driver有的是Java.sql.SQLException: The server time zone value有的是连接超时。原因驱动类名写错、驱动没被打进jar包、URL缺时区参数、数据库不可达这四类问题经常叠加出现让人误以为是Flink本身的问题。解决先在本地Java程序里用同样的URL和驱动类测试能否连上把Flink从我视野里摘出去能连上说明Flink侧问题连不上说明驱动和URL问题。确认后再检查jar包里的依赖用jar tf xxx.jar | grep mysql如果连驱动类都没有就是shade插件漏打了JdbcSink所在模块的依赖。5.4 自定义Source没有节流背压满红、CPU被打满现象作业“跑起来了”但TaskManager的CPU接近100%Web UI里背压页面一片HIGH数据产出量巨大但Sink根本写不过来。原因run方法里的while循环没有做任何节流collect频率远大于下游处理能力。很多人在测试环境用无限循环产生数据忘了加Thread.sleep导致背压从Sink一路传导到Source。解决在collect之后加一个可控的等待时间比如上述代码里的Thread.sleep(500)。如果想把速率做成可配置的可以从Configuration里读取参数。背压页面翻红不是程序崩溃但长期满背压会让checkpoint超时、状态发散这在生产里是大事。5.5 并行度比slot多作业一直等不到资源现象提交作业后Web UI显示作业状态为RUNNING但任务始终不调度时间戳一直不往前走日志里反复出现“Not enough tasks to saturate all slots”或类似信息。原因作业的并行度大于集群可用slot总数。本地默认模式只有1个TaskManager和1个slot你把并行度设成4还需要3个空余slot系统只能等。解决要么减少并行度用-p 1或-p 2把作业跑起来要么启动多个TaskManager。实验8阶段不要用修改conf里slot数量的方式解决那个影响面太大先学会准确判断资源再谈水平扩展。6. 用Web UI和REST API给作业做个体检并行度、背压和火焰图6.1 打开Web UI先看哪三样并行度、Back Pressure、火焰图作业能跑不代表作业正常这是实验8之后必须建立的直觉。每次提交作业后我习惯在localhost:8081上做三件事看Job Graph里每个算子的并行度是否如预期点进Back Pressure页面看背压状态出现LOW或HIGH就要追原因打开TaskManager页面找火焰图入口确认线程热点不在某个序列化或循环里。火焰图是Flink 1.14以后Web UI自带的能力它对排查CPU类问题非常直观。看到某个方法栈又宽又高通常就是性能瓶颈所在。很多背压问题看起来像玄学其实都是数据速率和Sink能力不匹配用火焰图把热点定位到具体类比瞎调并行度靠谱得多。6.2 用REST API把作业指标拽下来适合脚本化的体检方式Web UI适合人看但如果你想在每次提交后自动验证作业是否在正常处理数据可以走REST API# 拿到所有作业的概览 curl -s http://localhost:8081/jobs/overview | python3 -m json.tool # 记住你的jobid然后拉取实时速率指标 curl -s http://localhost:8081/jobs/{jobid}/metrics?getnumRecordsInPerSecond,numRecordsOutPerSecond两个指标分别是每秒进入数据流的记录数和每秒输出到Sink的记录数。如果numRecordsIn是零说明Source没有产出如果numRecordsIn很高但numRecordsOut是零说明中间算子卡死或Sink在抛异常。把这两条curl写进一个脚本每次提交后等十几秒再拉一次就能判断作业是不是真的在运转。带新人时我反复强调同一个习惯作业能跑不等于作业正常先看指标再改代码别急着加逻辑。这个习惯我从实验8一直用到现在几乎帮我躲过所有“看起来没报错但就是不处理数据”的雷。希望帮到你。本文还有配套的精品资源点击获取