Flink 快速入门
本页介绍 Flink 与 Hudi 的集成,并演示 Flink 如何为 Hudi 带来流处理的能力。
环境准备
Flink 支持矩阵
| Hudi | 支持的 Flink 版本 |
|---|---|
| 1.2.x | 1.17.x, 1.18.x, 1.19.x, 1.20.x(默认构建版本), 2.0.x, 2.1.x |
| 1.1.x | 1.17.x, 1.18.x, 1.19.x, 1.20.x(默认构建版本), 2.0.x |
| 1.0.x | 1.14.x, 1.15.x, 1.16.x, 1.17.x, 1.18.x, 1.19.x, 1.20.x(默认构建版本) |
| 0.15.x | 1.14.x, 1.15.x, 1.16.x, 1.17.x, 1.18.x |
| 0.14.x | 1.13.x, 1.14.x, 1.15.x, 1.16.x, 1.17.x |
下载 Flink 并启动 Flink 集群
- 你可以按照此处的说明来安装 Flink
- 然后在 Hadoop 环境中启动一个独立的 Flink 集群
对于本地环境,你可以下载 Hadoop 二进制包并如下设置 HADOOP_HOME:
# HADOOP_HOME is your hadoop root directory after unpack the binary package.
export HADOOP_CLASSPATH=`$HADOOP_HOME/bin/hadoop classpath`
# Start the Flink standalone cluster
./bin/start-cluster.sh请注意以下几点:
- 我们推荐使用 Hadoop 2.9.x 及以上版本,因为部分对象存储系统的文件系统实现仅从该版本起才提供
- flink-parquet 和 flink-flink 格式已打包进 hudi-flink-bundle jar 中
- Flink SQL
- DataStream API
我们使用 Flink SQL Client,因为它对 SQL 用户而言是一款优秀的快速上手工具。
启动 Flink SQL Client
Hudi 为 Fudi 提供了一个打包好的 bundle jar,在 Flink SQL Client 启动时需要加载该 jar。你可以在 hudi-source-dir/packaging/hudi-flink-bundle 路径下手动构建该 jar(参见 Build Flink Bundle Jar),或者从 Apache 官方仓库 下载。
现在启动 SQL CLI:
# Supported Flink versions for Hudi 1.2.x: 1.17, 1.18, 1.19, 1.20 (default build), 2.0, 2.1
export FLINK_VERSION=1.20
export HUDI_VERSION=1.2.0
wget https://repo1.maven.org/maven2/org/apache/hudi/hudi-flink${FLINK_VERSION}-bundle/${HUDI_VERSION}/hudi-flink${FLINK_VERSION}-bundle-${HUDI_VERSION}.jar -P /tmp/
./bin/sql-client.sh embedded -j /tmp/hudi-flink${FLINK_VERSION}-bundle-${HUDI_VERSION}.jar shell设置表名、基础路径,并通过 SQL 进行操作来开始本指南。SQL CLI 一次只能执行一行 SQL。
创建表
首先,我们来创建一个 Hudi 表。这里以分区表为例进行说明,但 Hudi 也支持非分区表。
- Flink SQL
- DataStream API
下面是一个创建 Flink Hudi 表的示例。
-- sets up the result mode to tableau to show the results directly in the CLI
set sql-client.execution.result-mode = tableau;
DROP TABLE hudi_table;
CREATE TABLE hudi_table(
ts BIGINT,
uuid VARCHAR(40) PRIMARY KEY NOT ENFORCED,
rider VARCHAR(20),
driver VARCHAR(20),
fare DOUBLE,
city VARCHAR(20)
)
PARTITIONED BY (`city`)
WITH (
'connector' = 'hudi',
'path' = 'file:///tmp/hudi_table',
'table.type' = 'MERGE_ON_READ'
);插入数据
- Flink SQL
- DataStream API
使用 SQL VALUES 向 Hudi 表中插入数据。
-- insert data using values
INSERT INTO hudi_table
VALUES
(1695159649087,'334e26e9-8355-45cc-97c6-c31daf0df330','rider-A','driver-K',19.10,'san_francisco'),
(1695091554788,'e96c4396-3fad-413a-a942-4cb36106d721','rider-C','driver-M',27.70 ,'san_francisco'),
(1695046462179,'9909a8b1-2d15-4d3d-8ec9-efc48c536a00','rider-D','driver-L',33.90 ,'san_francisco'),
(1695332066204,'1dced545-862b-4ceb-8b43-d2a568f6616b','rider-E','driver-O',93.50,'san_francisco'),
(1695516137016,'e3cf430c-889d-4015-bc98-59bdce1e530c','rider-F','driver-P',34.15,'sao_paulo'),
(1695376420876,'7a84095f-737f-40bc-b62f-6b69664712d2','rider-G','driver-Q',43.40 ,'sao_paulo'),
(1695173887231,'3eeb61f7-c2b0-4636-99bd-5d7a5a1d2c04','rider-I','driver-S',41.06 ,'chennai'),
(1695115999911,'c8abbe79-8d89-47ea-b4ce-4d224bae5bfa','rider-J','driver-T',17.85,'chennai');查询数据
- Flink SQL
- DataStream API
-- query from the Hudi table
select * from hudi_table;此语句查询数据集的快照视图。有关所有支持的表类型和查询类型的更多信息,请参阅表类型与查询。
更新数据
这与插入新数据类似。
- Flink SQL
- DataStream API
Hudi 表可以通过插入具有相同记录键(record key)的记录,或使用如下所示的标准 UPDATE 语句来进行更新。
-- Update Queries only works with batch execution mode
SET 'execution.runtime-mode' = 'batch';
UPDATE hudi_table SET fare = 25.0 WHERE uuid = '334e26e9-8355-45cc-97c6-c31daf0df330';note
UPDATE 语句自 Flink 1.17 起受支持,因此只有使用 Flink 1.17+ 编译的 Hudi Flink bundle 才提供此功能。只有在具有记录键(record key)的 Hudi 表上执行批查询时才能正确工作。
再次查询数据即可看到更新后的记录。每次写操作都会生成一个新的 commit,并以时间戳作为标识。
删除数据
- Flink SQL
- DataStream API
行级删除
在流式查询消费数据时,如果每行都设置了 RowKind,Hudi Flink source 也可以接收来自上游数据源的变更日志(change log),从而实现行级的 UPDATE 和 DELETE 操作。这样你就可以在 Hudi 上同步各类关系型数据库(RDBMS)的近实时快照。
批量删除
-- delete all the records with age greater than 23
-- NOTE: only works for batch sql queries
SET 'execution.runtime-mode' = 'batch';
DELETE FROM t1 WHERE age > 23;注意
DELETE语句自 Flink 1.17 起得到支持,因此只有使用 Flink 1.17 及以上版本编译的 Hudi Flink bundle 才提供该功能。只有对含有记录键(record key)的 Hudi 表进行批式查询时才能正确工作。
流式查询
Hudi Flink 还支持获取自某个给定提交时间戳以来发生变更的记录流。这可以通过使用 Hudi 的流式查询并提供一个变更需要开始流式传输的起始时间来实现。如果我们只需要给定提交之后的所有变更(通常正是这种情况),则无需指定 endTime。
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 is 60s.
);
-- Then query the table in stream mode
select * from t1;变更数据捕获查询
Hudi Flink 还提供了通过变更数据捕获(CDC)获取记录流的能力。CDC 查询适用于需要获取所有变更(包括记录变更前后的镜像)的应用场景。
set sql-client.execution.result-mode = tableau;
CREATE TABLE hudi_table(
ts BIGINT,
uuid VARCHAR(40) PRIMARY KEY NOT ENFORCED,
rider VARCHAR(20),
driver VARCHAR(20),
fare DOUBLE,
city VARCHAR(20)
)
PARTITIONED BY (`city`)
WITH (
'connector' = 'hudi',
'path' = 'file:///tmp/hudi_table',
'table.type' = 'COPY_ON_WRITE',
'cdc.enabled' = 'true' -- this option enables CDC logging
);
-- insert data using values
INSERT INTO hudi_table
VALUES
(1695159649087,'334e26e9-8355-45cc-97c6-c31daf0df330','rider-A','driver-K',19.10,'san_francisco'),
(1695091554788,'e96c4396-3fad-413a-a942-4cb36106d721','rider-C','driver-M',27.70 ,'san_francisco'),
(1695046462179,'9909a8b1-2d15-4d3d-8ec9-efc48c536a00','rider-D','driver-L',33.90 ,'san_francisco'),
(1695332066204,'1dced545-862b-4ceb-8b43-d2a568f6616b','rider-E','driver-O',93.50,'san_francisco'),
(1695516137016,'e3cf430c-889d-4015-bc98-59bdce1e530c','rider-F','driver-P',34.15,'sao_paulo'),
(1695376420876,'7a84095f-737f-40bc-b62f-6b69664712d2','rider-G','driver-Q',43.40 ,'sao_paulo'),
(1695173887231,'3eeb61f7-c2b0-4636-99bd-5d7a5a1d2c04','rider-I','driver-S',41.06 ,'chennai'),
(1695115999911,'c8abbe79-8d89-47ea-b4ce-4d224bae5bfa','rider-J','driver-T',17.85,'chennai');
SET 'execution.runtime-mode' = 'batch';
UPDATE hudi_table SET fare = 25.0 WHERE uuid = '334e26e9-8355-45cc-97c6-c31daf0df330';
-- Query the table in stream mode in another shell to see change logs
SET 'execution.runtime-mode' = 'streaming';
select * from hudi_table/*+ OPTIONS('read.streaming.enabled'='true')*/;这将返回 read.start-commit 提交之后发生的所有变更。该功能的独特之处在于,它让你能够在流式或批式数据源上编写流处理管道。
接下来该做什么?
- 快速入门:阅读上面的快速入门部分,快速上手使用 Flink SQL Client 向 Hudi 写入(并从 Hudi 读取)数据。
- 配置:全局配置通过
$FLINK_HOME/conf/flink-conf.yaml进行设置;作业级配置通过 Table Option 进行设置。 - 写入数据:Flink 支持多种写入模式,例如 CDC 摄取、批量插入、索引引导、Changelog 模式和 Append 模式。对于高吞吐量的追加管道,请选择 追加写入缓冲模式;对于大规模的更新插入(upsert)工作负载,请使用 记录级索引(RLI)分桶索引。Flink 还支持通过非阻塞并发控制实现多个流式写入器并发写入。
- 读取数据:Flink 支持多种读取模式,例如 流式查询和增量查询。若需更好的下推能力和可恢复的读取,请参阅 Flink Source V2;对于维度表关联,请使用 lookup join,并可选配堆外 RocksDB 缓存。
- 调优:针对写入/读取任务,本指南提供了一些调优建议,例如内存优化、托管内存写入缓冲以及写入速率限制。
- 优化:支持离线合并:离线合并。
- 查询引擎:除 Flink 外,还集成了许多其他引擎:Hive 查询、Presto 查询。
- Catalog:支持 Hudi 专属的 Catalog:Hudi Catalog。
如果你是 Apache Hudi 的新手,熟悉以下几个核心概念会很有帮助:
- Hudi 时间线 – Hudi 如何管理事务及其他表服务
- Hudi 存储布局 – 文件在存储上的组织方式
- Hudi 表类型 –
COPY_ON_WRITE和MERGE_ON_READ - Hudi 查询类型 – 快照查询、增量查询、读优化查询
更多内容请参见概念文档页面。
也可以浏览近期的博客文章,其中深入探讨了一些特定主题或用例。
Hudi 表可以通过 Hive、Spark、Flink、Presto 等多种查询引擎进行查询。我们制作了一个演示视频,在一个基于 Docker 的环境中展示了这些能力,所有依赖系统均在本地运行。我们建议你按照 docker demo 中的步骤复现相同的环境并亲自运行演示,以获得直观的体验。另外,如果你想了解如何将现有数据迁移到 Hudi,可参阅迁移指南。
评论
登录后参与评论
KnowForge