☰
Flink初级编程实践:从Socket WordCount到MySQL同步ClickHouse
2026/10/7 4:21:37 网站建设 项目流程

简介:这份资源是面向大数据课程学习者的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 11,JDK 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> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>1.18.1</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients</artifactId> <version>1.18.1</version> <scope>provided</scope> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.4.1</version> <executions> <execution> <phase>package</phase> <goals><goal>shade</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 是运行参数 DataStream<String> lines = env.socketTextStream("127.0.0.1", 9999); DataStream<Tuple2<String, Integer>> counts = lines .flatMap(new FlatMapFunction<String, Tuple2<String, Integer>>() { @Override public void flatMap(String line, Collector<Tuple2<String, 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 用nc,Windows 可以用 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 上定位算子。

DataStream<Tuple2<String, 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 初级编程实践中,很多人只设置了事件时间,忘了水位线,结果窗口一直不触发,作业看起来在跑但没有输出。

DataStream<Event> stream = env.addSource(new FlinkKafkaConsumer<>("topic", schema, props)) .assignTimestampsAndWatermarks( WatermarkStrategy.<Event>forBoundedOutOfOrderness(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 源 MySqlSource<String> mySqlSource = MySqlSource.<String>builder() .hostname("127.0.0.1") .port(3306) .databaseList("demo") // 监听的数据库 .tableList("demo.orders") // 监听的表 .username("root") .password("123456") .deserializer(new JsonDebeziumDeserializationSchema()) // 输出 JSON .build(); DataStream<String> 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)跑一遍全链路,确认数据能从源流到目标,再逐步调大并行度和批量参数。这样出问题时排查范围小,不用在集群日志里大海捞针。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询