快速入门

Flink 快速入门

师成师成· 更新于 2026-09-29· 阅读 23 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

本页介绍 Flink 与 Hudi 的集成,并演示 Flink 如何为 Hudi 带来流处理的能力。

环境准备

Flink 支持矩阵

Hudi支持的 Flink 版本
1.2.x1.17.x, 1.18.x, 1.19.x, 1.20.x(默认构建版本), 2.0.x, 2.1.x
1.1.x1.17.x, 1.18.x, 1.19.x, 1.20.x(默认构建版本), 2.0.x
1.0.x1.14.x, 1.15.x, 1.16.x, 1.17.x, 1.18.x, 1.19.x, 1.20.x(默认构建版本)
0.15.x1.14.x, 1.15.x, 1.16.x, 1.17.x, 1.18.x
0.14.x1.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 提交之后发生的所有变更。该功能的独特之处在于,它让你能够在流式或批式数据源上编写流处理管道。

接下来该做什么?

如果你是 Apache Hudi 的新手,熟悉以下几个核心概念会很有帮助:

更多内容请参见概念文档页面。

也可以浏览近期的博客文章,其中深入探讨了一些特定主题或用例。

Hudi 表可以通过 Hive、Spark、Flink、Presto 等多种查询引擎进行查询。我们制作了一个演示视频,在一个基于 Docker 的环境中展示了这些能力,所有依赖系统均在本地运行。我们建议你按照 docker demo 中的步骤复现相同的环境并亲自运行演示,以获得直观的体验。另外,如果你想了解如何将现有数据迁移到 Hudi,可参阅迁移指南。

评论

登录后参与评论

正在加载评论…