S1 第 1 期:不装 Flink 集群,一个 Java 程序跑通 Paimon
Apache Paimon 源码学习
作者 X老师(DaemonforY),Paimon 2.0 / master 源码,按 CC BY-NC-SA 4.0 发布。配套实验和代码在 GitHub。

TL;DR:Flink 的 Table API 自带一个嵌入式的 MiniCluster,不用下载和启动 Flink 集群,一个普通的 Java 程序就能执行 Flink SQL、读写 Paimon。好处是启动快、可以直接在 IDE 里对 Paimon 源码打断点。代价是要自己补齐几个依赖——补少了会遇到下面 4 个报错,而且有一个坑会让你多加一个没用的 jar。
为什么不用 Flink SQL Client

学 Paimon 最终要读源码、下断点。用 SQL Client 时,代码跑在独立的 Flink 进程里,调试要配远程参数、还要处理多个进程。而嵌入式运行时,SQL 解析、作业执行、Paimon 读写全在同一个 JVM 里,在 IDE 里点一下 Debug 就能停在 Paimon 的任意一行代码上。
第一步:一个最小的 Maven 工程
核心依赖只有这几类(版本以 Paimon 2.0.0 + Flink 2.2.0 为例):
| 依赖 | 作用 |
|---|---|
paimon-flink-2.2 | Paimon 的 Flink 连接器(含 Paimon 核心) |
flink-table-api-java-bridge、flink-table-planner_2.12、flink-table-runtime、flink-clients | 执行 Flink SQL |
程序本身只有几行:
TableEnvironment tEnv = TableEnvironment.create(EnvironmentSettings.inBatchMode());
tEnv.executeSql("CREATE CATALOG paimon WITH ('type'='paimon', 'warehouse'='file:///tmp/paimon')");
tEnv.executeSql("USE CATALOG paimon");
tEnv.executeSql("CREATE TABLE t (id BIGINT PRIMARY KEY NOT ENFORCED, v STRING)");
tEnv.executeSql("INSERT INTO t VALUES (1, 'a')").await();
tEnv.executeSql("SELECT * FROM t").print();看起来很简单。但第一次运行,我连续遇到了 4 个报错。
坑 1:ClassNotFoundException: org.apache.hadoop.conf.Configuration
现象:执行 INSERT 时,在 Paimon 的 FlinkSinkBuilder 里报错。 原因:Paimon 在 Flink 中需要 Hadoop 的类。生产环境中通常要把 Hadoop 的 jar 放进 flink/lib,嵌入式运行也一样。 解决:加 flink-shaded-hadoop-2-uber,并排除它的全部传递依赖(它是 uber jar,类都已经打在里面了,传递依赖反而会引入多余的日志实现)。
坑 2:NoClassDefFoundError: org/apache/log4j/Level
现象:作业提交后,在 Hadoop 的 UserGroupInformation 初始化时报错。 原因:Hadoop 用的是 log4j 1.x 的 API,而工程里没有任何 log4j。 解决:按 Flink 发行版的方式配齐日志:log4j-slf4j-impl、log4j-api、log4j-core,外加 log4j-1.2-api(把 log4j 1.x API 桥接到 log4j2)。顺便配一个 log4j2.properties,把 org.apache.paimon 设为 INFO——之后能看到 “Successfully commit snapshot 1” 这样的关键日志。
坑 3:NoClassDefFoundError: ...SingleThreadMultiplexSourceReaderBase

现象:写入成功了,但一 SELECT 就报错。 原因:Paimon 的读取端是标准的 FLIP-27 Source,它的 FileStoreSourceReader 继承自 org.apache.flink.connector.base 包下的这个类。 我的第一反应:类在 connector.base 包下,那就加 flink-connector-base——结果换来了坑 4。
坑 4:NoClassDefFoundError: ...BulkFormat$RecordIterator
现象:补了 flink-connector-base 后,SELECT 仍然报错,类名变了。 原因:Paimon 的 split reader 返回的迭代器类型来自 flink-connector-files。 解决:加 flink-connector-files。
然后我发现 flink-connector-base 是多余的:
$ unzip -l flink-connector-files-2.2.0.jar | grep -c org/apache/flink/connector/base/
108(108 是 jar 中该包路径下的条目数,含目录。Flink 1.20.1 的 flink-connector-files-1.20.1.jar 中是 107 个,结论相同。)
flink-connector-files 已经用 maven-shade-plugin 把 flink-connector-base 的类打包进来了。坑 3 和坑 4 其实是同一个依赖:一开始就只加 flink-connector-files,两个报错都不会出现。我在实验里把每个 jar 逐个去掉复现过:只去掉 connector-files → 坑 3;去掉它但补上 connector-base → 坑 4;只保留 connector-files、不要 connector-base → 读写全部成功。
额外的小坑:程序执行完了却不退出
MiniCluster 会启动一些非守护线程,main 方法结束(甚至抛异常)后 JVM 依然挂着。更隐蔽的是:如果你用管道过滤输出,异常信息被缓冲住,看起来像“卡死”。解决办法是在 main 结束时显式 System.exit()。
跑通之后
(节选,省略了 SQL 与结果之间的 INFO/WARN 日志)
Flink SQL> SELECT * FROM orders ORDER BY order_id;
+----------------------+----------------------+--------------+--------------------------------+--------------------------------+
| order_id | user_id | amount | status | dt |
+----------------------+----------------------+--------------+--------------------------------+--------------------------------+
| 1 | 101 | 99.90 | CREATED | 2026-09-24 |
| 2 | 102 | 15.00 | CREATED | 2026-09-24 |
| 3 | 103 | 250.00 | CREATED | 2026-09-24 |
| 4 | 101 | 8.80 | CREATED | 2026-09-25 |
| 5 | 104 | 66.60 | CREATED | 2026-09-25 |
+----------------------+----------------------+--------------+--------------------------------+--------------------------------+
5 rows in set写入 orders 时,日志里还能看到这一行:
20:54:14.267 INFO [FileStoreCommitImpl] Successfully commit snapshot 1 to table orders by user eda2c146-bb92-4071-8aa1-b47950558eec with identifier 9223372036854775807 and kind APPEND.9223372036854775807 是 Long.MAX_VALUE——批作业没有 checkpoint,输入结束时用它作为提交标识。这个数字后面讲 Flink 两阶段提交时还会再见到。
两个使用提醒
- JDK 17 需要加
--add-opens参数(Flink 在 JDK 17 上运行的要求;此说法基于 Flink 1.20.1 核对,2.2.0 未核对——本次 2.2.0 实验沿用了同一组参数),完整参数见仓库里的run.sh。 - 调试时把并行度设为 1、超时调大:停在断点上时,MiniCluster 可能因为心跳超时判定 TaskManager 失联。我在调试模式里把
heartbeat.timeout、pekko.ask.timeout设成了 1 小时。
小结
| 报错 | 出现在 | 缺的依赖 | Flink 发行版 lib/ 自带?(基于 Flink 1.20.1 核对,2.2.0 未核对) |
|---|---|---|---|
org.apache.hadoop.conf.Configuration | 写入 | flink-shaded-hadoop-2-uber(排除传递依赖) | ❌ 不带,生产上也要自己提供 |
org/apache/log4j/Level | 写入 | log4j2 全家桶 + log4j-1.2-api | ✅ |
SingleThreadMultiplexSourceReaderBase | 读取 | flink-connector-files | ✅ |
BulkFormat$RecordIterator | 读取 | flink-connector-files(同上) | ✅ |
4 个报错,只缺 3 个 jar。 规律:log4j 和 connector-files 是 Flink 发行版 lib/ 本来就有的;Hadoop 发行版不带(发行版内容基于 Flink 1.20.1 核对,2.2.0 未核对)——无论嵌入式还是集群部署,Paimon 都需要你自己提供 Hadoop(放入 shaded Hadoop jar 或设置 HADOOP_CLASSPATH)。
下期预告:表建好了,打开它的目录看看——snapshot、manifest、bucket,每个文件都是干什么的?
自测题:为什么写入成功、读取却失败?读和写分别用到了 Flink 的哪套接口?