快速入门
快速开始
提示
有关在 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 shellFlink 的 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 类型:
| Flink | Iceberg | 备注 |
|---|---|---|
| boolean | boolean | |
| tinyint | integer | |
| smallint | integer | |
| integer | integer | |
| bigint | long | |
| float | float | |
| double | double | |
| char | string | |
| varchar | string | |
| string | string | |
| binary | binary | |
| varbinary | fixed | |
| decimal | decimal | |
| date | date | |
| time | time | |
| timestamp | timestamp without timezone | |
| timestamp_ltz | timestamp with timezone | |
| array | list | |
| map | map | |
| multiset | map | |
| row | struct | |
| raw | 不支持 | |
| interval | 不支持 | |
| structured | 不支持 | |
| timestamp with zone | 不支持 | |
| distinct | 不支持 | |
| null | 不支持 | |
| symbol | 不支持 | |
| logical | 不支持 |
Iceberg 到 Flink
Iceberg 类型按下表转换为 Flink 类型:
| Iceberg | Flink | 备注 |
|---|---|---|
| boolean | boolean | |
| struct | row | |
| list | array | |
| map | map | |
| integer | integer | |
| long | bigint | |
| float | float | |
| double | double | |
| date | date | |
| time | time | |
| 不带时区的时间戳 | timestamp(6) | |
| 带时区的时间戳 | timestamp_ltz(6) | |
| string | varchar(2147483647) | |
| uuid | binary(16) | |
| fixed(N) | binary(N) | |
| binary | varbinary(2147483647) | |
| decimal(P, S) | decimal(P, S) | |
| 纳秒级时间戳 | timestamp(9) | |
| 纳秒级带时区时间戳 | timestamp_ltz(9) | |
| unknown | null | |
| variant | 不支持 | |
| geometry | 不支持 | |
| geography | 不支持 |
未来的改进
当前的 Flink Iceberg 集成尚不支持以下一些功能:
- 创建带隐藏分区的 Iceberg 表。Flink 邮件列表中的讨论。
- 创建带计算列的 Iceberg 表。
评论
登录后参与评论
KnowForge