尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Flink初级编程实践:从环境搭建到WordCount跑通与避坑指南

Flink初级编程实践:从环境搭建到WordCount跑通与避坑指南 简介面向大数据初学者的《实验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 UI6123是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里看到作业记录和运行时长。此时再跑一次无界流版本观察作业状态变成Runningdocker 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端口作为无界数据源 DataStreamString text env.socketTextStream(localhost, 9000); // 切割、计数、聚合 DataStreamTuple2String, Integer counts text .flatMap((String line, org.apache.flink.util.CollectorTuple2String, 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 cprint输出却是(c,1)先出现。原因很简单keyBy会把相同key路由到同一个子任务不同key分到不同子任务而print算子同样有并行度多个并行的print线程各自往stdout写谁先抢到输出缓冲区谁就显示在前面。这不是代码写错了是流计算的正常行为。想要保证输出顺序唯一的办法是把全局并行度设成1或者给keyBy之后的所有算子单独设置并行度1。我一般会在实验报告里写清这一点“并行度影响吞吐也影响输出顺序调试时优先串行压测时再调并行度。”面试时这也是高频追问点能说出这层关系基本就算理解了。4. 自定义数据源与数据汇实验里最值钱的部分4.1 自定义Data SourceSourceFunction与运行周期Flink自带的source类型有限fromElements、socketTextStream、fromFile都能用于实验但想模拟带频率的真实数据就要自己写SourceFunction。初级实验里最常用的模板是传感器数据源import org.apache.flink.streaming.api.functions.source.SourceFunction; public class SensorSource implements SourceFunctionString { private volatile boolean running true; private int counter 0; Override public void run(SourceContextString 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 runningcancel里置falserun里的wile循环自然会退出。不要靠Thread.stop()之类的手段Flink会标记作业失败而不是正常取消。第二ctx.collect()是唯一合法的输出方式不要直接在run里往外部写数据否则Flink的算子链和checkpoint机制全部失效。实际调试时我习惯在Source里保留休眠时间参数。Thread.sleep(1000)表示每秒一条想测试窗口计算就改成Thread.sleep(100)模拟高频数据想观察背压就把间隔调大到5秒看下游算子的处理节奏。这个参数在实验报告里能写出两组对比数据老师会很认可。4.2 自定义Data SinkJDBC连接器与必调参数实验要求把结果“落库”时最简单稳定的方案不是自定义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 ?; SinkFunctionTuple2String, 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都合理生产环境要压测后定。另一个是withDriverNameMySQL 8之前用com.mysql.jdbc.DriverMySQL 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写成了providedflink-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依赖里写了scopeprovided/scope这个scope的意思是“运行环境提供这个依赖”本地IDE自动把依赖加进classpath所以能跑提交时Flink集群没有这些类就炸了。解决将flink-streaming-java保持provided没问题因为集群确实自带Flink核心依赖但第三方库比如JSON库、Kafka客户端必须去掉provided。用maven-shade生成fat jar时排查mvn dependency:tree里有没有遗漏。这是一条最经典的生产级坑实验阶段踩一次比面试背十遍都管用。针对这个坑给出一个可直接套用的shade插件配置plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.5.1/version executions execution phasepackage/phase goalsgoalshade/goal/goals configuration transformers transformer implementationorg.apache.maven.plugins.shade.resource.ManifestResourceTransformer mainClasscom.example.SocketWordCount/mainClass /transformer /transformers /configuration /execution /executions /plugin5.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”骗过去。希望帮到你。本文还有配套的精品资源点击获取
返回列表