用 Flink 向 Paimon 写入持续更新的词频

用 Flink 向 Paimon 写入持续更新的词频

作者及维护方:Apache Paimon 项目/Apache Software Foundation。本文依据 Paimon 2.0《Quick Start》完整整理翻译,保留操作顺序、两种目录方式、内存配置和动态选项;版本与风险说明为编者补充。来源核对日期:2026-10-05。以下命令和 SQL 仅经静态审查,未在本次制作中执行。

这个例子把一条不断产生单字符词语的数据流,聚合成一张不断更新的词频表。它同时展示 Paimon 的两种读法:批查询取得已提交的表状态,流查询继续跟踪表的变化。数据库里的每个词对应一个主键,新的计数更新这个词的记录,而不是把每一次累计结果都当作互不相关的新行。

数据生成器把词送入 Flink 分组计数,成功检查点提交 Paimon 主键表,批查询读取快照,流查询跟踪变化
词频写入与两种读取路径。编者依据官方教程自绘,非运行截图。

先配对 Flink 与 Paimon 连接器

原页列出 Flink 2.2、2.1、2.0、1.20、1.19、1.18、1.17 和 1.16 对应的 Paimon 2.0.0 bundled JAR。它建议使用较新的 Flink 版本,但安装时仍必须选对连接器的 Flink 版本段。用于读写的 bundled JAR 与用于手动压缩等维护操作的 action JAR 是两类文件,不能互相替代。

以 Flink 1.20 系列为例,目标文件是 paimon-flink-1.20-2.0.0.jar;action 文件是 paimon-flink-action-2.0.0.jar。本次下载并静态核对了 Maven Central 上已发布的 paimon-flink-1.20:2.0.0 POM:Flink 依赖版本为 1.20.1,作用域为 provided,运行环境仍需提供 Flink。另保存了 paimon-hadoop-uber:2.0.0 的发布 POM。它们提供依赖边界证据,不等同于 JAR 已通过运行兼容性测试;本次没有加载或执行这些组件。

也可以从 官方源码仓库构建。原文命令如下:

mvn clean install -DskipTests

构建产物分别位于 ./paimon-flink/paimon-flink-<flink-version>/target/paimon-flink-<flink-version>-2.0.0.jar 和 ./paimon-flink/paimon-flink-action/target/paimon-flink-action-2.0.0.jar。clean 会清理构建输出,-DskipTests 会跳过测试;构建成功并不等于测试通过。源码目录应检出与目标版本对应的发布标签,实际文件名以该检出的构建结果为准。

准备本地集群

先从 Flink 官方下载页取得目标发行包并解压。原文用 tar -xzf flink-*.tgz,再用通配符复制连接器。为避免同目录多个版本被同时选中,下面把连接器固定为 1.20 对应文件;/path/to/flink 必须替换为你的实际目录。这是对原文通配符的明确收紧,不是已执行过的安装记录。

tar -xzf flink-1.20.1-bin-scala_2.12.tgz
cp paimon-flink-1.20-2.0.0.jar /path/to/flink/lib/

解压会写入当前目录,复制 JAR 会改变集群类路径。不要把多个 Flink 版本的 Paimon 连接器一起放入 lib。上面的发行包文件名是本例所选版本的说明,下载渠道、校验值、Java 要求应按该发行版说明核对。

原文第三步要求准备 Hadoop 依赖:若机器已有 Hadoop 环境,应保证 HADOOP_CLASSPATH 包含常用 Hadoop 库;否则原文给出把 flink-shaded-hadoop-2-uber-*.jar 放入 Flink lib 的旧式命令。这个通配文件名和指向下载页的链接不足以确认任何当前可用、适配的 JAR。Paimon 2.0 文件系统说明明确将本地文件系统列为内建支持,并为缺少 Hadoop 依赖、又需要访问 Hadoop 集群的应用提供 paimon-hadoop-uber-2.0.0.jar。本例使用本地 file: 仓库,不应为照抄旧步骤而盲目加入重复 Hadoop 依赖;需要 HDFS 时,再按同版文档补齐类路径与 Hadoop 配置。

为了同时运行写入作业与查询作业,集群至少需要两个任务槽。Flink 1.19 以前编辑 conf/flink-conf.yaml,1.19 及以后使用 conf/config.yaml。原文的配置项是:

taskmanager.numberOfTaskSlots: 2

这是本地教程的最低资源安排,不是任意并行度都足够的容量配置。在 Linux/类 Unix Shell 中启动集群和 SQL 客户端:

/path/to/flink/bin/start-cluster.sh
/path/to/flink/bin/sql-client.sh

在浏览器访问 http://localhost:8081,检查 Flink 控制台中的集群状态。启动脚本会创建本地服务和监听端口;此处的控制台只用于本机演示,不应未经访问控制直接暴露到公网。Windows 环境不能把这些 Bash 脚本直接当作 PowerShell 命令。

创建目录与主键表

普通 Paimon Catalog 使用仓库路径保存表文件:

CREATE CATALOG my_catalog WITH (
    'type' = 'paimon',
    'warehouse' = 'file:/tmp/paimon'
);
USE CATALOG my_catalog;

CREATE TABLE word_count (
    word STRING PRIMARY KEY NOT ENFORCED,
    cnt BIGINT
);

word 是词,cnt 是累计次数。PRIMARY KEY NOT ENFORCED 是 Flink SQL 的主键声明,不能理解为数据库将替你检查所有输入数据的唯一性。这里的聚合按词生成更新,Paimon 按主键维护相应记录。

file:/tmp/paimon 仅适用于这台机器上的演示;/tmp 也不是长期数据保留承诺。分布式运行时,各节点必须能访问同一个仓库,应换成 HDFS、OSS 等共享文件系统,并按其要求配置插件和权限。不要让多个节点各自把同名本地目录当作共享仓库。

原文还给出另一条路径:使用依赖 Hive Metastore 的 FlinkGenericCatalog,在同一个目录里管理 Paimon、Hive 以及 Kafka 等 Flink 通用表。此时创建 Paimon 表要显式指定连接器:

CREATE CATALOG my_catalog WITH (
    'type' = 'paimon-generic',
    'hive-conf-dir' = '/path/to/hive/conf',
    'hadoop-conf-dir' = '/path/to/hadoop/conf'
);
USE CATALOG my_catalog;

CREATE TABLE word_count (
    word STRING PRIMARY KEY NOT ENFORCED,
    cnt BIGINT
) WITH (
    'connector' = 'paimon'
);

这段与前一个目录方案二选一,不能在同一会话里直接重复创建同名目录。路径由原文的 ... 改成了明显的占位路径。Paimon 会使用 hive-site.xml 中的 hive.metastore.warehouse.dir;应带上 hdfs:// 等 URI scheme,否则可能被解释为本地路径。

持续生成词并写入累计结果

临时表不把生成的数据作为持久源表保存,它通过 Flink 的 datagen 连接器不断产生长度为 1 的字符串。随后按词分组计数,把动态结果写入 Paimon:

CREATE TEMPORARY TABLE word_table (
    word STRING
) WITH (
    'connector' = 'datagen',
    'fields.word.length' = '1'
);

SET 'execution.checkpointing.interval' = '10 s';

INSERT INTO word_count
SELECT word, COUNT(*) FROM word_table GROUP BY word;

流式写入需要配置检查点。这里的 10 s 是检查点触发间隔,不是“每十秒一定能读到新数据”的服务保证。只有写入作业成功完成相应检查点、数据提交成功,读者才能看到相应提交。作业调度、背压、存储响应和失败恢复都会影响可见时间。若表暂时为空,应先看作业与检查点状态,而不是反复重建表。

用批查询观察已提交状态

保持写入作业运行,调整 SQL 客户端的查询配置:

SET 'sql-client.execution.result-mode' = 'tableau';
RESET 'execution.checkpointing.interval';
SET 'execution.runtime-mode' = 'batch';

SELECT * FROM word_count;

tableau 让结果在终端中按表格展示。切换到批模式之后,这次查询读取表的有限状态并结束。多执行几次,可以观察写入提交后词频如何变化。这里重置的是会话中后续作业使用的检查点选项,不是在远程取消已经提交的流式写入作业。

切回流模式,追踪计数区间的变化

第二种读取方式让查询持续运行:

SET 'execution.runtime-mode' = 'streaming';

SELECT `interval`, COUNT(*) AS interval_cnt
FROM (
    SELECT cnt / 10000 AS `interval`
    FROM word_count
)
GROUP BY `interval`;

内层把累计词频按 cnt / 10000 归类,外层统计每个结果值对应多少个词。当某个词的累计计数变化、跨过分组边界时,查询的分组统计也会更新。这里的 interval 是列别名,反引号用来避免与 SQL 关键字冲突;它不是时间窗口。这个演示使用整数计数,具体表达式的类型与除法语义仍应以所选 Flink 版本的 SQL 类型推导为准。

流式结果可能表现为插入、更新或撤回形式的变化,不能把每一条终端输出都当成新的独立事实累加。批模式回答“当前已提交表是什么样”,流模式回答“表接下来发生了什么变化”。本次未运行集群,因此没有附造出来的终端表格或性能数字。

退出时同时处理作业、客户端和集群

先在 Flink 控制台取消持续运行的写入作业和流式查询,再退出 SQL 客户端:

-- 仅在确实要删除表及其数据时,才自行取消下一行注释。
-- DROP TABLE word_count;

EXIT;

DROP TABLE 具有删除表及清理数据文件的风险,原文也将它留在注释中;本文保持这一默认。退出客户端不等于所有已经提交的作业都结束。最后停止本地集群:

/path/to/flink/bin/stop-cluster.sh

让 Flink 管理写入缓冲内存

教程最后介绍 Paimon sink 使用执行器内存池的方式。默认情况下,每个任务自行分配和管理堆内存缓冲;当一个 TaskManager 上的任务很多时,这些独立内存池可能造成性能问题,甚至 OOM。启用 Flink managed memory 后,由 Flink 协调 writer buffer 的资源分配。

选项 默认值 作用
sink.use-managed-memory-allocator false 设为 true 时,Flink sink 的 merge tree 使用托管内存;否则使用各任务独立的分配器。
sink.managed.writer-buffer-memory 256M writer buffer 在托管内存分配中的权重/请求量;实际可用资源取决于运行环境与 Flink 的分配。

原文第二个选项的说明先强调按权重分配、实际值依赖环境,随后又称当前该值等于运行时实际分配量,表述存在张力。不能据此向任意集群承诺每个 writer 都恰好获得 256 MiB;应结合目标版本实现与运行指标核对。

INSERT INTO paimon_table
/*+ OPTIONS(
    'sink.use-managed-memory-allocator' = 'true',
    'sink.managed.writer-buffer-memory' = '256M'
) */
SELECT * FROM source_table;

原文用 SELECT * FROM .... 表示尚未填写的查询;这里改成了命名占位表 source_table,读者仍需先创建真实源表和目标表。它说明的是 hint 的位置,并非可以直接接着前面执行的完整作业。

用动态选项覆盖当前作业的表配置

有些参数可以临时调整,而不修改 Catalog 中保存的表选项。表级动态键格式为 paimon.${catalogName}.${dbName}.${tableName}.${config_key};目录、数据库或表名位置可以用 * 通配。全局动态键直接使用 ${config_key},对所有表生效。出现冲突时,表级选项覆盖全局选项。

-- 所有表使用这个时间戳
SET 'scan.timestamp-millis' = '1697018249001';
SELECT * FROM T;

-- 指定目录、数据库和表
SET 'paimon.mycatalog.default.T.scan.timestamp-millis' = '1697018249000';
SELECT * FROM T;

-- 任意目录的 default.T
SET 'paimon.*.default.T.scan.timestamp-millis' = '1697018249000';
SELECT * FROM T;

-- T1 与其他表采用不同的扫描时间戳
SET 'paimon.mycatalog.default.T1.scan.timestamp-millis' = '1697018249000';
SET 'scan.timestamp-millis' = '1697018249001';
-- SELECT * FROM T1 JOIN T2 ON <实际连接条件>;

这些时间戳是原文演示值,表名也是示意。它们不保证你的仓库保留了相应快照,更不适合无条件复制进当前词频演示。最后一行原文写成 ON xxxx,本文改为注释中的占位条件,避免把非完整 SQL 伪装成可直接运行的代码。动态选项在当前会话的作业配置中生效;试验结束应清理不再需要的覆盖项,避免后续查询误用历史时间点。

来源:Quick Start(Paimon 2.0);补充核对:Filesystems。版权:© 2023–2026 The Apache Software Foundation;Apache Paimon、Paimon 及其标志是 ASF 商标。Apache Paimon 项目按 Apache License 2.0 提供代码;本译稿为中文翻译整理,保留项目版权与来源,并在下方附完整许可文本及与本文相关的 NOTICE 署名。新增说明与原创配图均已标明。

Apache License 2.0

Apache Paimon 代码采用此许可;完整许可文本如下。

                                 Apache License
                           Version 2.0, January 2004
                        http://www.apache.org/licenses/

   TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION

   1. Definitions.

      "License" shall mean the terms and conditions for use, reproduction,
      and distribution as defined by Sections 1 through 9 of this document.

      "Licensor" shall mean the copyright owner or entity authorized by
      the copyright owner that is granting the License.

      "Legal Entity" shall mean the union of the acting entity and all
      other entities that control, are controlled by, or are under common
      control with that entity. For the purposes of this definition,
      "control" means (i) the power, direct or indirect, to cause the
      direction or management of such entity, whether by contract or
      otherwise, or (ii) ownership of fifty percent (50%) or more of the
      outstanding shares, or (iii) beneficial ownership of such entity.

      "You" (or "Your") shall mean an individual or Legal Entity
      exercising permissions granted by this License.

      "Source" form shall mean the preferred form for making modifications,
      including but not limited to software source code, documentation
      source, and configuration files.

      "Object" form shall mean any form resulting from mechanical
      transformation or translation of a Source form, including but
      not limited to compiled object code, generated documentation,
      and conversions to other media types.

      "Work" shall mean the work of authorship, whether in Source or
      Object form, made available under the License, as indicated by a
      copyright notice that is included in or attached to the work
      (an example is provided in the Appendix below).

      "Derivative Works" shall mean any work, whether in Source or Object
      form, that is based on (or derived from) the Work and for which the
      editorial revisions, annotations, elaborations, or other modifications
      represent, as a whole, an original work of authorship. For the purposes
      of this License, Derivative Works shall not include works that remain
      separable from, or merely link (or bind by name) to the interfaces of,
      the Work and Derivative Works thereof.

      "Contribution" shall mean any work of authorship, including
      the original version of the Work and any modifications or additions
      to that Work or Derivative Works thereof, that is intentionally
      submitted to Licensor for inclusion in the Work by the copyright owner
      or by an individual or Legal Entity authorized to submit on behalf of
      the copyright owner. For the purposes of this definition, "submitted"
      means any form of electronic, verbal, or written communication sent
      to the Licensor or its representatives, including but not limited to
      communication on electronic mailing lists, source code control systems,
      and issue tracking systems that are managed by, or on behalf of, the
      Licensor for the purpose of discussing and improving the Work, but
      excluding communication that is conspicuously marked or otherwise
      designated in writing by the copyright owner as "Not a Contribution."

      "Contributor" shall mean Licensor and any individual or Legal Entity
      on behalf of whom a Contribution has been received by Licensor and
      subsequently incorporated within the Work.

   2. Grant of Copyright License. Subject to the terms and conditions of
      this License, each Contributor hereby grants to You a perpetual,
      worldwide, non-exclusive, no-charge, royalty-free, irrevocable
      copyright license to reproduce, prepare Derivative Works of,
      publicly display, publicly perform, sublicense, and distribute the
      Work and such Derivative Works in Source or Object form.

   3. Grant of Patent License. Subject to the terms and conditions of
      this License, each Contributor hereby grants to You a perpetual,
      worldwide, non-exclusive, no-charge, royalty-free, irrevocable
      (except as stated in this section) patent license to make, have made,
      use, offer to sell, sell, import, and otherwise transfer the Work,
      where such license applies only to those patent claims licensable
      by such Contributor that are necessarily infringed by their
      Contribution(s) alone or by combination of their Contribution(s)
      with the Work to which such Contribution(s) was submitted. If You
      institute patent litigation against any entity (including a
      cross-claim or counterclaim in a lawsuit) alleging that the Work
      or a Contribution incorporated within the Work constitutes direct
      or contributory patent infringement, then any patent licenses
      granted to You under this License for that Work shall terminate
      as of the date such litigation is filed.

   4. Redistribution. You may reproduce and distribute copies of the
      Work or Derivative Works thereof in any medium, with or without
      modifications, and in Source or Object form, provided that You
      meet the following conditions:

      (a) You must give any other recipients of the Work or
          Derivative Works a copy of this License; and

      (b) You must cause any modified files to carry prominent notices
          stating that You changed the files; and

      (c) You must retain, in the Source form of any Derivative Works
          that You distribute, all copyright, patent, trademark, and
          attribution notices from the Source form of the Work,
          excluding those notices that do not pertain to any part of
          the Derivative Works; and

      (d) If the Work includes a "NOTICE" text file as part of its
          distribution, then any Derivative Works that You distribute must
          include a readable copy of the attribution notices contained
          within such NOTICE file, excluding those notices that do not
          pertain to any part of the Derivative Works, in at least one
          of the following places: within a NOTICE text file distributed
          as part of the Derivative Works; within the Source form or
          documentation, if provided along with the Derivative Works; or,
          within a display generated by the Derivative Works, if and
          wherever such third-party notices normally appear. The contents
          of the NOTICE file are for informational purposes only and
          do not modify the License. You may add Your own attribution
          notices within Derivative Works that You distribute, alongside
          or as an addendum to the NOTICE text from the Work, provided
          that such additional attribution notices cannot be construed
          as modifying the License.

      You may add Your own copyright statement to Your modifications and
      may provide additional or different license terms and conditions
      for use, reproduction, or distribution of Your modifications, or
      for any such Derivative Works as a whole, provided Your use,
      reproduction, and distribution of the Work otherwise complies with
      the conditions stated in this License.

   5. Submission of Contributions. Unless You explicitly state otherwise,
      any Contribution intentionally submitted for inclusion in the Work
      by You to the Licensor shall be under the terms and conditions of
      this License, without any additional terms or conditions.
      Notwithstanding the above, nothing herein shall supersede or modify
      the terms of any separate license agreement you may have executed
      with Licensor regarding such Contributions.

   6. Trademarks. This License does not grant permission to use the trade
      names, trademarks, service marks, or product names of the Licensor,
      except as required for reasonable and customary use in describing the
      origin of the Work and reproducing the content of the NOTICE file.

   7. Disclaimer of Warranty. Unless required by applicable law or
      agreed to in writing, Licensor provides the Work (and each
      Contributor provides its Contributions) on an "AS IS" BASIS,
      WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
      implied, including, without limitation, any warranties or conditions
      of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
      PARTICULAR PURPOSE. You are solely responsible for determining the
      appropriateness of using or redistributing the Work and assume any
      risks associated with Your exercise of permissions under this License.

   8. Limitation of Liability. In no event and under no legal theory,
      whether in tort (including negligence), contract, or otherwise,
      unless required by applicable law (such as deliberate and grossly
      negligent acts) or agreed to in writing, shall any Contributor be
      liable to You for damages, including any direct, indirect, special,
      incidental, or consequential damages of any character arising as a
      result of this License or out of the use or inability to use the
      Work (including but not limited to damages for loss of goodwill,
      work stoppage, computer failure or malfunction, or any and all
      other commercial damages or losses), even if such Contributor
      has been advised of the possibility of such damages.

   9. Accepting Warranty or Additional Liability. While redistributing
      the Work or Derivative Works thereof, You may choose to offer,
      and charge a fee for, acceptance of support, warranty, indemnity,
      or other liability obligations and/or rights consistent with this
      License. However, in accepting such obligations, You may act only
      on Your own behalf and on Your sole responsibility, not on behalf
      of any other Contributor, and only if You agree to indemnify,
      defend, and hold each Contributor harmless for any liability
      incurred by, or claims asserted against, such Contributor by reason
      of your accepting any such warranty or additional liability.

   END OF TERMS AND CONDITIONS

   APPENDIX: How to apply the Apache License to your work.

      To apply the Apache License to your work, attach the following
      boilerplate notice, with the fields enclosed by brackets "[]"
      replaced with your own identifying information. (Don't include
      the brackets!)  The text should be enclosed in the appropriate
      comment syntax for the file format. We also recommend that a
      file or class name and description of purpose be included on the
      same "printed page" as the copyright notice for easier
      identification within third-party archives.

   Copyright [yyyy] [name of copyright owner]

   Licensed under the Apache License, Version 2.0 (the "License");
   you may not use this file except in compliance with the License.
   You may obtain a copy of the License at

       http://www.apache.org/licenses/LICENSE-2.0

   Unless required by applicable law or agreed to in writing, software
   distributed under the License is distributed on an "AS IS" BASIS,
   WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
   See the License for the specific language governing permissions and
   limitations under the License.

Apache Paimon NOTICE

Apache Paimon
Copyright 2023-2026 The Apache Software Foundation

This product includes software developed at
The Apache Software Foundation (http://www.apache.org/).

相关上游文件:Apache Paimon 2.0.0 LICENSE、Apache Paimon 2.0.0 NOTICE。附录保留本稿内容所涉及的项目归属信息。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容