Flink SQL 查询引擎连接器
Apache Paimon(孵化中)
登录后可跨设备保存划线和私人笔记登录
Apache Paimon (Incubating) 是一个流式数据湖平台,支持高速数据摄入、变更数据跟踪以及高效的实时分析。
Tip
本文假设你已掌握 Apache Paimon (Incubating) 的基础知识与操作。对于本文未涉及的知识,你可以从其官方文档中获取。
通过使用 Kyuubi,我们可以对 Apache Paimon (Incubating) 运行 SQL 查询,这种方式比直接使用 Flink 更加方便、易于理解且易于扩展。
Apache Paimon (Incubating) 集成
要启用 Kyuubi Flink SQL 引擎与 Apache Paimon (Incubating) 的集成,你需要:
- 引入 Apache Paimon (Incubating) 的依赖项
依赖项
支持 Apache Paimon (Incubating) 的 Kyuubi Flink SQL 引擎的 classpath 由以下部分组成:
- kyuubi-flink-sql-engine-1.9.1_2.12.jar,即随 Kyuubi 发行版部署的引擎 jar;
- 一份 Flink 发行版的副本;
- paimon-flink-<version>.jar(例如:paimon-flink-1.16-0.4-SNAPSHOT.jar),可在 Apache Paimon (Incubating) 支持的引擎 Flink 中找到;
- flink-shaded-hadoop-2-uber-<version>.jar,其代码可在预打包的 Hadoop Jar 中找到。
为了让 Apache Paimon (Incubating) 的包对引擎的运行时 classpath 可见,你需要:
- 将 Apache Paimon (Incubating) 的包直接放入
$FLINK_HOME/lib; - 设置 HADOOP_CLASSPATH 环境变量,或将预打包的 Hadoop Jar 复制到 flink/lib。
Warning
请注意不同 Apache Paimon (Incubating) 与 Flink 版本之间的兼容性,可以在 Apache Paimon (Incubating) 多引擎支持页面上确认。
Apache Paimon (Incubating) 操作
以 CREATE CATALOG 为例,
CREATE CATALOG my_catalog WITH (
'type'='paimon',
'warehouse'='file:/tmp/paimon'
);
USE CATALOG my_catalog;以 CREATE TABLE 为例,
CREATE TABLE MyTable (
user_id BIGINT,
item_id BIGINT,
behavior STRING,
dt STRING,
PRIMARY KEY (dt, user_id) NOT ENFORCED
) PARTITIONED BY (dt) WITH (
'bucket' = '4'
);以 Query Table 为例,
SET 'execution.runtime-mode' = 'batch';
SELECT * FROM orders WHERE catalog_id=1025;以 Streaming Query 为例,
SET 'execution.runtime-mode' = 'streaming';
SELECT * FROM MyTable /*+ OPTIONS ('log.scan'='latest') */;以 Rescale Bucket 为例,
ALTER TABLE my_table SET ('bucket' = '4');
INSERT OVERWRITE my_table PARTITION (dt = '2022-01-01');评论
登录后参与评论
正在加载评论…
KnowForge