Flink SQL 查询引擎连接器
Hudi
登录后可跨设备保存划线和私人笔记登录
Apache Hudi(读作“hoodie”)是新一代流式数据湖平台。Apache Hudi 将核心的数据仓库和数据库功能直接带入数据湖。
提示
本文假定你已经掌握了 Hudi 的基础知识和操作。对于本文中未提及的 Hudi 相关知识,你可以查阅其官方文档。
通过使用 Kyuubi,我们可以对 Hudi 执行 SQL 查询,这比直接使用 Flink 操作 Hudi 更加方便、易懂且易于扩展。
Hudi 集成
要通过 Catalog API 启用 Kyuubi Flink SQL 引擎与 Hudi 的集成,你需要:
- 引入 Hudi 的依赖
依赖
支持 Hudi 的 Kyuubi Flink SQL 引擎的 classpath 由以下部分组成:
- kyuubi-flink-sql-engine-1.9.1_2.12.jar,即随 Kyuubi 发行版一起部署的引擎 jar 包;
- 一份 Flink 发行版的拷贝;
- hudi-flink<flink.version>-bundle_<scala.version>-<hudi.version>.jar(示例:hudi-flink1.14-bundle_2.12-0.11.1.jar),可在 Maven Central 中找到。
为了让 Hudi 包对引擎的运行时 classpath 可见,我们可以使用以下方法之一:
- 将 Hudi 包直接放入
$flink_HOME/lib; - 设置
pipeline.jars=/path/to/hudi-flink-bundle。
Hudi 操作
以 Create Table 为例:
CREATE TABLE t1 (
id INT PRIMARY KEY NOT ENFORCED,
name STRING,
price DOUBLE
) WITH (
'connector' = 'hudi',
'path' = 's3://bucket-name/hudi/',
'table.type' = 'MERGE_ON_READ' -- this creates a MERGE_ON_READ table, by default is COPY_ON_WRITE
);以 Query Data 为例,
SELECT * FROM t1;以 Insert and Update Data 为例,
INSERT INTO t1 VALUES (1, 'Lucas' , 2.71828);以 Streaming Query 为例,
CREATE TABLE t1 (
uuid VARCHAR(20) PRIMARY KEY NOT ENFORCED,
name VARCHAR(10),
age INT,
ts TIMESTAMP(3),
`partition` VARCHAR(20)
)
PARTITIONED BY (`partition`)
WITH (
'connector' = 'hudi',
'path' = '${path}',
'table.type' = 'MERGE_ON_READ',
'read.streaming.enabled' = 'true', -- this option enable the streaming read
'read.start-commit' = '20210316134557', -- specifies the start commit instant time
'read.streaming.check-interval' = '4' -- specifies the check interval for finding new source commits, default 60s.
);
-- Then query the table in stream mode
SELECT * FROM t1;以 Delete Data 为例:
流式查询可以隐式地自动删除数据。在流式查询消费数据时,Hudi Flink source 也可以接收底层数据源的变更日志,进而按行级别应用 UPDATE 和 DELETE 操作。
评论
登录后参与评论
正在加载评论…
KnowForge