Skip to content

S1 第 1 期:不装 Flink 集群,一个 Java 程序跑通 Paimon ​

Apache Paimon 源码学习

作者 X老师(DaemonforY),Paimon 2.0 / master 源码,按 CC BY-NC-SA 4.0 发布。配套实验和代码在 GitHub。

Yui和Kai用嵌入式程序读写Paimon
Yui和Kai用嵌入式程序读写PaimonAI 生成配图

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

同一JVM内调试执行并读写Paimon
同一JVM内调试执行并读写PaimonAI 生成配图

学 Paimon 最终要读源码、下断点。用 SQL Client 时,代码跑在独立的 Flink 进程里,调试要配远程参数、还要处理多个进程。而嵌入式运行时,SQL 解析、作业执行、Paimon 读写全在同一个 JVM 里,在 IDE 里点一下 Debug 就能停在 Paimon 的任意一行代码上。

第一步:一个最小的 Maven 工程 ​

核心依赖只有这几类(版本以 Paimon 2.0.0 + Flink 2.2.0 为例):

依赖作用
paimon-flink-2.2Paimon 的 Flink 连接器(含 Paimon 核心)
flink-table-api-java-bridge、flink-table-planner_2.12、flink-table-runtime、flink-clients执行 Flink SQL

程序本身只有几行:

java
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 ​

Paimon读取端通过连接桥接入文件读取器
Paimon读取端通过连接桥接入文件读取器AI 生成配图

现象:写入成功了,但一 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 两阶段提交时还会再见到。

两个使用提醒 ​

  1. JDK 17 需要加 --add-opens 参数(Flink 在 JDK 17 上运行的要求;此说法基于 Flink 1.20.1 核对,2.2.0 未核对——本次 2.2.0 实验沿用了同一组参数),完整参数见仓库里的 run.sh。
  2. 调试时把并行度设为 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 的哪套接口?

代码示例在页面里运行时使用 HiveGPT 的模型接口。延伸阅读来自 JavaGuide(Apache-2.0),版权归原作者。