
简介这份资源是面向大数据课程学习者的Flink初级编程实验报告对应「大数据技术原理与应用」课程实验8适合正在学习流处理框架、需要完成实验作业或复盘操作流程的高校学生与初学者。压缩包内仅含1个docx文档约2.46MB以图文并茂的实验报告形式呈现完整记录实验环境、操作截图与结果验证。内容围绕两个核心任务展开一是使用IntelliJ IDEA开发WordCount程序涵盖Flink与Maven安装、Java代码编写、打包JAR并提交集群运行二是借助Linux自带的NC程序模拟实时数据流编写Flink程序完成词频统计并部署运行。报告还整理了Idea引用Flink报错、Maven打包缓慢、NC程序无输出等常见问题的排查与解决思路并附有Flink Web控制台查看输出的方法。目前已有5153人学习下载可帮助读者快速掌握Flink基本开发流程、环境搭建与调试技巧。1. Flink初级编程实践从一条无界流到可复现的本地作业很多人第一次接触 Flink 初级编程实践卡住的地方不是算子写不出来而是环境跑不通、依赖对不上、作业提交后看不到输出。我带过几批新人最常见的场景是照着示例写了一个从 Socket 读数据、做 WordCount 的作业mvn package成功flink run却报ClassNotFoundException或者作业在 Web UI 上显示 RUNNING但一条结果都不打印。这篇笔记就围绕这个真实痛点展开把 Flink 初级编程实践拆成「环境怎么搭、DataStream API 怎么写、参数怎么调、坑在哪」四段让你能在本地把一条无界流跑通再迁移到真实数据源。适合刚上手 Flink、需要交一份能跑的作业、或者准备把 MySQL 同步到 ClickHouse 这类链路做原型验证的工程师。下面所有命令和代码都以本地单机 Standalone 模式为基准不依赖任何云服务。2. 环境搭建与第一个 DataStream 作业把依赖和入口类先钉死2.1 版本对齐Flink、Scala、JDK 三者的绑定关系Flink 初级编程实践翻车最多的地方就是版本。Flink 1.18 之后默认不再捆绑 Scala如果你用 Scala API必须自己引入flink-scala对应版本用 Java API 则相对干净。我一般建议新手先用 Java把算子逻辑跑通再考虑 Scala。JDK 方面Flink 1.15 及以上推荐 JDK 11JDK 8 在部分连接器上会有UnsupportedClassVersionError。下面这张表是我实际项目里验证过的组合直接抄即可。组件推荐版本说明JDK11JDK 17 在部分连接器上仍有反射限制Flink1.18.1稳定版社区文档齐全Maven3.8低于 3.6 会出现依赖解析异常Scala2.12仅在使用 Scala API 时需要提示不要混用 Flink 1.14 的代码和 1.18 的运行时StreamExecutionEnvironment的部分方法签名已经变化编译能过但运行会抛NoSuchMethodError。2.2 Maven 依赖与打包插件的最小配置依赖写不对作业提交必挂。核心是flink-streaming-java和flink-clients前者提供 DataStream API后者提供本地执行环境。打包时用maven-shade-plugin把业务类打进去Flink 自身依赖设为provided避免和集群里的包冲突。dependencies !-- Flink 核心依赖集群已提供打包时排除 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version1.18.1/version scopeprovided/scope /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-clients/artifactId version1.18.1/version scopeprovided/scope /dependency /dependencies build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.4.1/version executions execution phasepackage/phase goalsgoalshade/goal/goals /execution /executions /plugin /plugins /build逻辑说明provided表示编译时需要、打包时不带入因为 Flink 集群的lib目录已经有这些 jar。如果你用 IDE 本地直接 run需要把provided临时改成compile否则会报类找不到。参数上flink-streaming-java的版本必须和集群lib里的版本完全一致差一个小版本都可能出问题。2.3 第一个 WordCount从 Socket 读无界流下面这段代码是 Flink 初级编程实践的标准入口从本地 9999 端口读文本按空格切词5 秒滚动窗口统计。它覆盖了 Source、Transformation、Sink 三个环节是理解 DataStream 模型的最小闭环。public class SocketWordCount { public static void main(String[] args) throws Exception { // 创建本地执行环境并行度设为 1 便于观察输出 StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); // 从 socket 读取无界流host 和 port 是运行参数 DataStreamString lines env.socketTextStream(127.0.0.1, 9999); DataStreamTuple2String, Integer counts lines .flatMap(new FlatMapFunctionString, Tuple2String, Integer() { Override public void flatMap(String line, CollectorTuple2String, Integer out) { for (String word : line.split(\\s)) { if (!word.isEmpty()) { out.collect(Tuple2.of(word, 1)); } } } }) .keyBy(value - value.f0) // 按单词分组 .window(TumblingProcessingTimeWindows.of(Time.seconds(5))) // 5 秒滚动窗口 .sum(1); // 对第二个字段求和 counts.print(); // 输出到标准输出 env.execute(Socket WordCount); // 触发执行 } }逻辑说明socketTextStream产生的是一个无界流flatMap把每行拆成单词并转成(word,1)二元组keyBy按单词做逻辑分区window定义 5 秒的滚动窗口sum(1)对计数累加。参数上setParallelism(1)是为了让print()的输出顺序可读生产环境按 CPU 核数设置。TumblingProcessingTimeWindows用的是处理时间不依赖事件时间戳适合入门如果要处理乱序数据需要换成TumblingEventTimeWindows并配合水位线。2.4 本地运行与提交命令先在终端启动一个 TCP 服务端再运行作业。Linux 和 macOS 用ncWindows 可以用 PowerShell 的Test-NetConnection或直接写一个 Java Socket Server。# 终端 1启动本地 socket 服务监听 9999 nc -lk 9999 # 终端 2编译打包 mvn clean package -DskipTests # 终端 2本地提交作业到 Standalone 集群 ./bin/flink run -c com.example.SocketWordCount target/flink-demo-1.0.jar逻辑说明nc -lk 9999中的-l表示监听-k表示保持连接这样 Flink 作业不会因为服务端断开而退出。flink run的-c指定入口类全限定名jar 路径必须是 shade 之后的包。提交后在 Web UI 的Task Managers里能看到print算子的输出或者在提交作业的终端直接看到计数结果。如果看不到输出先检查setParallelism是否大于 1 导致输出分散再检查 socket 是否真的在监听。3. 常用算子与窗口参数把 keyBy、窗口、水位线讲透3.1 keyBy 的分区逻辑与并行度陷阱keyBy是 Flink 初级编程实践里最容易误解的算子。它不是把数据按 key 分组后放到一个算子实例里而是按key.hashCode() % 并行度做逻辑分区相同 key 一定进同一个 subtask。这意味着如果你把并行度从 1 改成 4同一个单词仍然只会在一个 subtask 里累加但不同单词会分散到不同 subtask。很多人看到print输出里同一个单词出现在多个 subtask就以为 keyBy 失效了其实是并行度设置和输出顺序的问题。参数上keyBy的 key 必须可序列化且不能为 null否则会抛NullPointerException。如果 key 是自定义对象建议实现hashCode和equals否则分区结果不稳定。我一般会在 keyBy 之后加一个uid方便在 Web UI 上定位算子。DataStreamTuple2String, Integer counts lines .flatMap(...) .keyBy(value - value.f0) .window(TumblingProcessingTimeWindows.of(Time.seconds(5))) .sum(1) .uid(word-count-window); // 固定算子 ID便于状态恢复逻辑说明uid是算子的唯一标识Flink 保存点依赖它来恢复状态。如果不设置Flink 会自动生成一个但代码改动后可能变化导致保存点无法恢复。参数上uid一旦上线就不要改否则恢复会失败。3.2 窗口类型选择滚动、滑动、会话的适用场景窗口是 Flink 初级编程实践的核心概念选错窗口类型会导致结果不符合预期。滚动窗口按固定时长切分窗口之间不重叠滑动窗口有滑动步长窗口之间可能重叠会话窗口按活动间隙切分适合用户行为分析。下面这张表对比了三种窗口的关键参数。窗口类型构造参数适用场景注意点滚动窗口Time.seconds(5)固定周期统计窗口边界对齐处理时间滑动窗口Time.seconds(10), Time.seconds(5)平滑统计窗口重叠导致重复计算会话窗口Time.seconds(30)用户会话分析间隙内无数据则关闭窗口注意处理时间窗口的结果依赖数据到达顺序事件时间窗口需要配合水位线否则窗口永远不触发。3.3 水位线与事件时间乱序数据的后悔药当数据源带事件时间戳且可能乱序时必须设置水位线。水位线是一个时间戳表示「早于这个时间的数据已经全部到达」。Flink 初级编程实践中很多人只设置了事件时间忘了水位线结果窗口一直不触发作业看起来在跑但没有输出。DataStreamEvent stream env.addSource(new FlinkKafkaConsumer(topic, schema, props)) .assignTimestampsAndWatermarks( WatermarkStrategy.EventforBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) - event.getEventTime()) ); stream.keyBy(Event::getUserId) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .sum(amount) .print();逻辑说明forBoundedOutOfOrderness(Duration.ofSeconds(5))表示允许 5 秒的乱序水位线会延迟 5 秒推进。withTimestampAssigner从事件里提取时间戳。参数上乱序容忍度要根据实际数据延迟设置设太小会丢迟到数据设太大会增加结果延迟。如果数据延迟超过容忍度可以配置allowedLateness让窗口保留一段时间接收迟到数据。4. 避坑与排查Flink 初级编程实践里最常见的 5 个翻车现场4.1 作业提交报 ClassNotFoundException现象flink run提交后立刻失败日志里出现ClassNotFoundException: com.example.SocketWordCount。原因打包时没有把业务类打进去或者-c指定的类名写错。解决检查maven-shade-plugin是否生效用jar tf target/xxx.jar | grep SocketWordCount确认类在包里-c后面的类名必须是全限定名区分大小写。4.2 作业 RUNNING 但 print 无输出现象Web UI 显示作业 RUNNING但提交终端和 TaskManager 日志都没有计数结果。原因print算子的输出被并行度分散或者 socket 源没有数据流入。解决先把并行度设为 1确认 socket 服务端有数据再检查print是否被setParallelism覆盖。如果用的是flink run -d后台提交输出不会回到终端需要去 TaskManager 的.out文件里看。4.3 窗口不触发结果一直不输出现象作业运行正常但窗口统计结果迟迟不打印。原因用了事件时间窗口但没设置水位线或者水位线推进太慢。解决确认assignTimestampsAndWatermarks已调用且forBoundedOutOfOrderness的容忍度不要设得过大如果数据源本身没有事件时间改用处理时间窗口。4.4 状态后端配置错误导致 OOM现象作业运行一段时间后 TaskManager 内存溢出日志出现OutOfMemoryError。原因默认状态后端把状态放在 JVM 堆内存窗口多、key 多时堆内存不够。解决在flink-conf.yaml里配置 RocksDB 状态后端把状态放到磁盘。state.backend: rocksdb state.backend.incremental: true state.checkpoints.dir: hdfs:///flink/checkpoints逻辑说明rocksdb把状态存储在本地磁盘适合大状态场景incremental: true开启增量检查点减少每次 checkpoint 的数据量。参数上state.checkpoints.dir需要提前创建目录否则 checkpoint 会失败。4.5 依赖冲突导致 NoSuchMethodError现象作业提交后抛NoSuchMethodError或NoClassDefFoundError。原因业务 jar 里打入了和集群版本不一致的 Flink 依赖或者引入了冲突的第三方库。解决用mvn dependency:tree检查依赖树把 Flink 相关依赖设为provided如果必须引入第三方库用 shade 插件的relocation重命名包路径避免和集群冲突。5. 从本地到真实链路把 MySQL 同步到 ClickHouse 的进阶技巧Flink 初级编程实践跑通 WordCount 之后下一步通常是接真实数据源。热搜里「使用 flink 实现 mysql 同步到 clickhouse」是很多人的目标这里给一个可落地的思路。核心是用 Flink CDC 连接器读 MySQL 的 binlog经过简单转换后写入 ClickHouse。注意这不是初级作业的必选项但能帮你验证这套编程模型在真实链路里的边界。第一步引入 Flink CDC 和 ClickHouse JDBC 连接器。CDC 连接器版本要和 Flink 版本对齐比如 Flink 1.18 对应 CDC 3.0 左右。第二步用MySqlSource构建源表开启增量快照。第三步用JdbcSink写入 ClickHouse注意批量提交参数。// 构建 MySQL CDC 源 MySqlSourceString mySqlSource MySqlSource.Stringbuilder() .hostname(127.0.0.1) .port(3306) .databaseList(demo) // 监听的数据库 .tableList(demo.orders) // 监听的表 .username(root) .password(123456) .deserializer(new JsonDebeziumDeserializationSchema()) // 输出 JSON .build(); DataStreamString source env.fromSource( mySqlSource, WatermarkStrategy.noWatermarks(), MySQL Source ); // 写入 ClickHouse批量 100 条或 1 秒提交一次 source.addSink(JdbcSink.sink( INSERT INTO orders (id, amount) VALUES (?, ?), (statement, record) - { JSONObject json JSON.parseObject(record); statement.setInt(1, json.getInteger(id)); statement.setBigDecimal(2, json.getBigDecimal(amount)); }, JdbcExecutionOptions.builder() .withBatchSize(100) .withBatchIntervalMs(1000) .withMaxRetries(3) .build(), new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() .withUrl(jdbc:clickhouse://127.0.0.1:8123/demo) .withDriverName(com.clickhouse.jdbc.ClickHouseDriver) .build() ));逻辑说明MySqlSource的databaseList和tableList支持正则JsonDebeziumDeserializationSchema把变更事件转成 JSON 字符串包含op字段标识增删改。JdbcSink的withBatchSize控制批量提交条数withBatchIntervalMs控制提交间隔两者满足其一就触发。参数上ClickHouse 的 JDBC URL 默认端口是 8123驱动类名根据你用的驱动包调整。如果写入报Too many parts把批量调大或间隔调长。验证方法在 MySQL 里执行一条INSERT观察 ClickHouse 里是否出现对应记录再执行UPDATE和DELETE确认 CDC 能捕获变更。如果只同步了全量没有增量检查 MySQL 的binlog_format是否为ROW以及用户是否有REPLICATION SLAVE权限。我自己的习惯是任何 Flink 作业上线前先在本地用setParallelism(1)跑一遍全链路确认数据能从源流到目标再逐步调大并行度和批量参数。这样出问题时排查范围小不用在集群日志里大海捞针。希望帮到你。本文还有配套的精品资源点击获取