Apache Flink

快速入门

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

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

快速开始

提示

有关在 Flink 中使用 Iceberg 的概述,请参阅 Flink 快速入门。

Apache Iceberg 同时支持 Apache Flink 的 DataStream API 和 Table API。有关 Apache Flink 的集成信息,请参阅多引擎支持页面。

功能支持Flink备注
SQL 创建 catalog✔️
SQL 创建 database✔️
SQL 创建表✔️
SQL create table like✔️
SQL 修改表✔️仅支持修改表属性,不支持列和分区的变更
SQL 删除表✔️
SQL 查询✔️同时支持流式和批处理模式
SQL insert into✔️ ️同时支持流式和批处理模式
SQL insert overwrite✔️ ️
DataStream 读取✔️ ️
DataStream 追加✔️ ️
DataStream 覆写✔️ ️
元数据表✔️
重写文件操作✔️ ️

使用 Flink SQL Client 的准备工作

要在 Flink 中创建 Iceberg 表,建议使用 Flink SQL Client,因为这样更容易让用户理解相关概念。

从 Apache 下载页面 下载 Flink。Iceberg 在编译 Apache iceberg-flink-runtime jar 时使用的是 Scala 2.12,因此建议使用随 Scala 2.12 打包的 Flink 2.3。

FLINK_VERSION=2.3.0
SCALA_VERSION=2.12
APACHE_FLINK_URL=https://archive.apache.org/dist/flink/
wget ${APACHE_FLINK_URL}/flink-${FLINK_VERSION}/flink-${FLINK_VERSION}-bin-scala_${SCALA_VERSION}.tgz
tar xzvf flink-${FLINK_VERSION}-bin-scala_${SCALA_VERSION}.tgz

在 Hadoop 环境中启动独立的 Flink 集群:

# HADOOP_HOME is your hadoop root directory after unpack the binary package.
APACHE_HADOOP_URL=https://archive.apache.org/dist/hadoop/
HADOOP_VERSION=2.8.5
wget ${APACHE_HADOOP_URL}/common/hadoop-${HADOOP_VERSION}/hadoop-${HADOOP_VERSION}.tar.gz
tar xzvf hadoop-${HADOOP_VERSION}.tar.gz
HADOOP_HOME=`pwd`/hadoop-${HADOOP_VERSION}

export HADOOP_CLASSPATH=`$HADOOP_HOME/bin/hadoop classpath`

# Start the flink standalone cluster
cd flink-${FLINK_VERSION}/
./bin/start-cluster.sh

启动 Flink SQL 客户端。Iceberg 项目中有一个独立的 flink-runtime 模块,用于生成一个打包 jar,可由 Flink SQL 客户端直接加载。要手动构建 flink-runtime 打包 jar,请构建 iceberg 项目,构建后会在 <iceberg-root-dir>/flink-runtime/build/libs 下生成该 jar。或者,从 Apache 仓库 下载 flink-runtime jar。

# HADOOP_HOME is your hadoop root directory after unpack the binary package.
export HADOOP_CLASSPATH=`$HADOOP_HOME/bin/hadoop classpath`

# Below works for Flink 1.15 or earlier
./bin/sql-client.sh embedded -j <flink-runtime-directory>/iceberg-flink-runtime-2.3-1.11.0.jar shell

# Flink 1.16+ has a regression in loading external jars via -j. See FLINK-30035 for details.
# put iceberg-flink-runtime-2.3-1.11.0.jar in flink/lib dir
./bin/sql-client.sh embedded shell

默认情况下,Iceberg 随附了用于 Hadoop catalog 的 Hadoop jar。若要使用 Hive catalog,请在打开 Flink SQL 客户端时加载 Hive jar。幸运的是,Flink 已经为 SQL 客户端提供了捆绑的 Hive jar。以下示例演示如何下载依赖并开始使用:

# HADOOP_HOME is your hadoop root directory after unpack the binary package.
export HADOOP_CLASSPATH=`$HADOOP_HOME/bin/hadoop classpath`

ICEBERG_VERSION=1.11.0
MAVEN_URL=https://repo1.maven.org/maven2
ICEBERG_MAVEN_URL=${MAVEN_URL}/org/apache/iceberg
ICEBERG_PACKAGE=iceberg-flink-runtime
FLINK_VERSION_MAJOR=2.3
wget ${ICEBERG_MAVEN_URL}/${ICEBERG_PACKAGE}-${FLINK_VERSION_MAJOR}/${ICEBERG_VERSION}/${ICEBERG_PACKAGE}-${FLINK_VERSION_MAJOR}-${ICEBERG_VERSION}.jar -P lib/

HIVE_VERSION=2.3.9
SCALA_VERSION=2.12
FLINK_VERSION=2.3.0
FLINK_CONNECTOR_URL=${MAVEN_URL}/org/apache/flink
FLINK_CONNECTOR_PACKAGE=flink-sql-connector-hive
wget ${FLINK_CONNECTOR_URL}/${FLINK_CONNECTOR_PACKAGE}-${HIVE_VERSION}_${SCALA_VERSION}/${FLINK_VERSION}/${FLINK_CONNECTOR_PACKAGE}-${HIVE_VERSION}_${SCALA_VERSION}-${FLINK_VERSION}.jar

./bin/sql-client.sh embedded shell

Flink 的 Python API

Info

PyFlink 1.6.1 在 Apple Silicon 的 macOS 上存在已知问题。参见 FLINK-28786。

使用 pip 安装 Apache Flink 依赖:

pip install apache-flink==2.3.0

提供一个指向 iceberg-flink-runtime jar 的 file:// 路径,该 jar 可以通过构建项目后在 <iceberg-root-dir>/flink-runtime/build/libs 目录下找到,或者从 Apache 官方仓库 下载。可以通过以下方式向 pyflink 添加第三方 jar:

  • env.add_jars("file:///my/jar/path/connector.jar")
  • table_env.get_config().get_configuration().set_string("pipeline.jars", "file:///my/jar/path/connector.jar")

官方文档中也提到了这一点。下面的示例使用 env.add_jars(..):

import os

from pyflink.datastream import StreamExecutionEnvironment

env = StreamExecutionEnvironment.get_execution_environment()
iceberg_flink_runtime_jar = os.path.join(os.getcwd(), "iceberg-flink-runtime-2.3-1.11.0.jar")

env.add_jars("file://{}".format(iceberg_flink_runtime_jar))

接下来,创建一个 StreamTableEnvironment 并执行 Flink SQL 语句。下面的示例展示了如何通过 Python Table API 创建自定义 catalog:

from pyflink.table import StreamTableEnvironment
table_env = StreamTableEnvironment.create(env)
table_env.execute_sql("""
CREATE CATALOG my_catalog WITH (
    'type'='iceberg',
    'catalog-impl'='com.my.custom.CatalogImpl',
    'my-additional-catalog-config'='my-value'
)
""")

运行查询:

(table_env
    .sql_query("SELECT PULocationID, DOLocationID, passenger_count FROM my_catalog.nyc.taxis LIMIT 5")
    .execute()
    .print())
+----+----------------------+----------------------+--------------------------------+
| op |         PULocationID |         DOLocationID |                passenger_count |
+----+----------------------+----------------------+--------------------------------+
| +I |                  249 |                   48 |                            1.0 |
| +I |                  132 |                  233 |                            1.0 |
| +I |                  164 |                  107 |                            1.0 |
| +I |                   90 |                  229 |                            1.0 |
| +I |                  137 |                  249 |                            1.0 |
+----+----------------------+----------------------+--------------------------------+
5 rows in set

更多详情请参阅 Python Table API。

添加 Catalog

Flink 支持使用 Flink SQL 创建 Catalog。

Catalog 配置

执行以下查询即可创建并命名一个 Catalog(请将 <catalog_name> 替换为你的 Catalog 名称,将 '<config_key>' = '<config_value>' 替换为 Catalog 实现的配置):

CREATE CATALOG <catalog_name> WITH (
  'type'='iceberg',
  '<config_key>' = '<config_value>'
);

以下属性可以全局设置,不受特定目录(catalog)实现的限制:

  • type:必须为 iceberg。(必填)
  • catalog-type:内置目录可选 hive、hadoop、rest、glue、jdbc 或 nessie;若使用 catalog-impl 自定义目录实现,则可不设置该属性。(可选)
  • catalog-impl:自定义目录实现的全限定类名。如果未设置 catalog-type,则必须设置该属性。(可选)
  • property-version:用于描述属性版本的版本号。当属性格式发生变化时,该属性可用于向后兼容。当前的属性版本为 1。(可选)
  • cache-enabled:是否启用目录缓存,默认值为 true。(可选)
  • cache.expiration-interval-ms:目录条目在本地缓存的时长,单位为毫秒;-1 等负值表示禁用过期,不允许设置为 0。默认值为 -1。(可选)

Hive 目录

这会创建一个名为 hive_catalog 的 Iceberg 目录,可以通过 'catalog-type'='hive' 进行配置,它从 Hive metastore 加载表:

CREATE CATALOG hive_catalog WITH (
  'type'='iceberg',
  'catalog-type'='hive',
  'uri'='thrift://localhost:9083',
  'clients'='5',
  'property-version'='1',
  'warehouse'='hdfs://nn:8020/warehouse/path'
);

REST catalog

这将创建一个名为 rest_catalog 的 Iceberg catalog,可通过 'catalog-type'='rest' 进行配置,它从 REST catalog 加载表:

CREATE CATALOG rest_catalog WITH (
  'type'='iceberg',
  'catalog-type'='rest',
  'uri'='https://localhost/'
);

创建表

CREATE TABLE `hive_catalog`.`default`.`sample` (
    id BIGINT COMMENT 'unique id',
    data STRING
);

写入

若要通过 Flink 流式作业向表中追加新数据,请使用 INSERT INTO:

INSERT INTO `hive_catalog`.`default`.`sample` VALUES (1, 'a');
INSERT INTO `hive_catalog`.`default`.`sample` SELECT id, data from other_kafka_table;

要在批处理作业中用查询结果替换表中的数据,请使用 INSERT OVERWRITE(Flink 流处理作业不支持 INSERT OVERWRITE)。对于 Iceberg 表,覆盖是原子操作。

SELECT 查询产生数据行的分区将被替换,例如:

INSERT OVERWRITE `hive_catalog`.`default`.`sample` VALUES (1, 'a');

Iceberg 还支持通过 SELECT 查询的值覆盖指定分区:

INSERT OVERWRITE `hive_catalog`.`default`.`sample` PARTITION(data='a') SELECT 6;

Flink 原生支持将 DataStream<RowData> 和 DataStream<Row> 写入 Iceberg 表 sink。

StreamExecutionEnvironment env = ...;

DataStream<RowData> input = ... ;
Configuration hadoopConf = new Configuration();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path", hadoopConf);

FlinkSink.forRowData(input)
    .tableLoader(tableLoader)
    .append();

env.execute("Test Iceberg DataStream");

分支写入

Iceberg 表同样支持通过 FlinkSink 中的 toBranch API 向分支写入数据。

关于分支的更多信息,请参阅分支。

FlinkSink.forRowData(input)
    .tableLoader(tableLoader)
    .toBranch("audit-branch")
    .append();

读取

使用以下语句提交一个 Flink 批处理作业:

-- Execute the flink job in batch mode for current session context
SET execution.runtime-mode = batch;
SELECT * FROM `hive_catalog`.`default`.`sample`;

Iceberg 支持在 流式 Flink 作业中处理增量数据,此类作业可从某个历史快照 ID 开始:

-- Submit the flink job in streaming mode for current session.
SET execution.runtime-mode = streaming;

-- Enable this switch because streaming read SQL will provide few job options in flink SQL hint options.
SET table.dynamic-table-options.enabled=true;

-- Read all the records from the iceberg current snapshot, and then read incremental data starting from that snapshot.
SELECT * FROM `hive_catalog`.`default`.`sample` /*+ OPTIONS('streaming'='true', 'monitor-interval'='1s')*/ ;

-- Read all incremental data starting from the snapshot-id '3821550127947089987' (records from this snapshot will be excluded).
SELECT * FROM `hive_catalog`.`default`.`sample` /*+ OPTIONS('streaming'='true', 'monitor-interval'='1s', 'start-snapshot-id'='3821550127947089987')*/ ;

SQL 也是检查表的推荐方式。要查看表中的所有快照,请使用 snapshots 元数据表:

SELECT * FROM `hive_catalog`.`default`.`sample$snapshots`;

Iceberg 的 Java API 支持流式读取或批量读取:

DataStream<RowData> batch = FlinkSource.forRowData()
     .env(env)
     .tableLoader(tableLoader)
     .streaming(false)
     .build();

类型转换

Iceberg 对 Flink 的集成会自动在 Flink 类型与 Iceberg 类型之间进行转换。当向某个表写入该表使用了 Flink 不支持的类型(例如 UUID)的数据时,Iceberg 会接收来自 Flink 类型的值并进行转换。

Flink 到 Iceberg

Flink 类型按下表转换为 Iceberg 类型:

FlinkIceberg备注
booleanboolean
tinyintinteger
smallintinteger
integerinteger
bigintlong
floatfloat
doubledouble
charstring
varcharstring
stringstring
binarybinary
varbinaryfixed
decimaldecimal
datedate
timetime
timestamptimestamp without timezone
timestamp_ltztimestamp with timezone
arraylist
mapmap
multisetmap
rowstruct
raw不支持
interval不支持
structured不支持
timestamp with zone不支持
distinct不支持
null不支持
symbol不支持
logical不支持

Iceberg 到 Flink

Iceberg 类型按下表转换为 Flink 类型:

IcebergFlink备注
booleanboolean
structrow
listarray
mapmap
integerinteger
longbigint
floatfloat
doubledouble
datedate
timetime
不带时区的时间戳timestamp(6)
带时区的时间戳timestamp_ltz(6)
stringvarchar(2147483647)
uuidbinary(16)
fixed(N)binary(N)
binaryvarbinary(2147483647)
decimal(P, S)decimal(P, S)
纳秒级时间戳timestamp(9)
纳秒级带时区时间戳timestamp_ltz(9)
unknownnull
variant不支持
geometry不支持
geography不支持

未来的改进

当前的 Flink Iceberg 集成尚不支持以下一些功能:

  • 创建带隐藏分区的 Iceberg 表。Flink 邮件列表中的讨论。
  • 创建带计算列的 Iceberg 表。

评论

登录后参与评论

正在加载评论…