探索Iceberg与FlinkSQL的深度集成之旅
Flink 与 Iceberg 集成指南
环境准备
安装 Flink
-
Flink 与 Iceberg 版本对应关系
Flink 版本 Iceberg 版本 1.11 0.9.0 – 0.12.1 1.12 0.12.0 – 0.13.1 1.13 0.13.0 – 1.0.0 1.14 0.13.0 – 1.1.0 1.15 0.14.0 – 1.1.0 1.16 1.1.0 – 1.1.0 -
上传并解压 Flink 安装包
bashtar -zxvf flink-1.16.0-bin-scala_2.12.tgz -C /opt/module/ -
配置环境变量
bashsudo vim /etc/profile.d/my_env.sh export HADOOP_CLASSPATH=hadoop classpath source /etc/profile.d/my_env.sh -
拷贝 Iceberg 的 jar 包到 Flink 的 lib 目录
bashcp /opt/software/iceberg/iceberg-flink-runtime-1.16-1.1.0.jar /opt/module/flink-1.16.0/lib
启动 Hadoop
启动 SQL-Client
-
修改
flink-conf.yaml配置yamlclassloader.check-leaked-classloader: false taskmanager.numberOfTaskSlots: 4 state.backend: rocksdb execution.checkpointing.interval: 30000 state.checkpoints.dir: hdfs://hadoop1:8020/ckps state.backend.incremental: true -
启动 Flink
bash/opt/module/flink-1.16.0/bin/start-cluster.sh -
启动 Flink 的 SQL-Client
bash/opt/module/flink-1.16.0/bin/sql-client.sh embedded
创建和使用 Catalog
语法说明
sql
CREATE CATALOG <catalog_name> WITH (
'type'='iceberg',
<config_key>=<config_value>
);
- type: 必须是
iceberg。(必填) - catalog-type: 内置了
hive和hadoop两种 catalog,也可以使用catalog-impl来自定义 catalog。(可选) - catalog-impl: 自定义 catalog 实现的全限定类名。如果未设置
catalog-type,则必须设置。(可选) - property-version: 描述属性版本的版本号。此属性可用于向后兼容,以防属性格式更改。当前属性版本为
1。(可选) - cache-enabled: 是否启用目录缓存,默认值为
true。(可选) - cache.expiration-interval-ms: 本地缓存 catalog 条目的时间(以毫秒为单位);负值,如
-1表示没有时间限制,不允许设为0。默认值为-1。(可选)
Hive Catalog
-
上传 Hive Connector 到 Flink 的 lib 中
bashcp flink-sql-connector-hive-3.1.2_2.12-1.16.0.jar /opt/module/flink-1.16.0/lib/ -
启动 Hive Metastore 服务
bashhive --service metastore -
创建 Hive Catalog
sqlCREATE CATALOG hive_catalog WITH ( 'type'='iceberg', 'catalog-type'='hive', 'uri'='thrift://192.168.110.120:9083', 'clients'='5', 'property-version'='1', 'warehouse'='hdfs://192.168.110.120:8020/warehouse/iceberg-hive' ); USE CATALOG hive_catalog;
Hadoop Catalog
sql
CREATE CATALOG hadoop_catalog WITH (
'type'='iceberg',
'catalog-type'='hadoop',
'warehouse'='hdfs://192.168.110.120:8020/warehouse/iceberg-hadoop',
'property-version'='1'
);
USE CATALOG hadoop_catalog;
配置 SQL-Client 初始化文件
sql
## hive catalog
CREATE CATALOG hive_catalog WITH (
'type'='iceberg',
'catalog-type'='hive',
'uri'='thrift://192.168.110.120:9083',
'clients'='5',
'property-version'='1',
'warehouse'='hdfs://192.168.110.120:8020/warehouse/iceberg-hive'
);
## hadoop catalog
CREATE CATALOG hadoop_catalog WITH (
'type'='iceberg',
'catalog-type'='hadoop',
'warehouse'='hdfs://192.168.110.120:8020/warehouse/iceberg-hadoop',
'property-version'='1'
);
## 默认使用 hive_catalog(可选)
USE CATALOG hive_catalog;
启动 SQL-Client 时,加上 -i 参数指定初始化文件:
bash
/opt/module/flink-1.16.0/bin/sql-client.sh embedded -i conf/sql-client-init.sql
DDL 语句
创建数据库
sql
CREATE DATABASE iceberg_db;
USE iceberg_db;
创建表
sql
CREATE TABLE hive_catalog.default.sample (
id BIGINT COMMENT 'unique id',
data STRING
);
创建分区表
sql
CREATE TABLE `hive_catalog`.`default`.`sample` (
id BIGINT COMMENT 'unique id',
data STRING
) PARTITIONED BY (data);
使用 LIKE 语法建表
sql
CREATE TABLE `hive_catalog`.`default`.`sample_like` LIKE `hive_catalog`.`default`.`sample`;
修改表
-
修改表属性
sqlALTER TABLE hive_catalog.default.sample SET ('write.format.default'='avro'); -
修改表名
sqlALTER TABLE hive_catalog.default.sample RENAME TO hive_catalog.default.new_sample;
删除表
sql
DROP TABLE hive_catalog.default.sample;
插入语句
INSERT INTO
sql
INSERT INTO hive_catalog.default.sample VALUES (1, 'a');
INSERT INTO hive_catalog.default.sample SELECT id, data FROM sample2;
INSERT OVERWRITE
sql
SET execution.runtime-mode = batch;
INSERT OVERWRITE sample VALUES (1, 'a');
INSERT OVERWRITE hive_catalog.default.sample PARTITION(data='a') SELECT 6;
UPSERT
-
建表时指定
sqlCREATE TABLE hive_catalog.iceberg_db.sample ( id INT UNIQUE COMMENT 'unique id', data STRING NOT NULL, PRIMARY KEY(id) NOT ENFORCED ) WITH ( 'format-version'='2', 'write.upsert.enabled'='true' ); -
插入时指定
sqlINSERT INTO tableName /*+ OPTIONS('upsert-enabled'='true') */ ... -
读取 Kafka 流,upsert 插入到 Iceberg 表中
sqlCREATE TABLE default_catalog.default_database.kafka ( id INT, data STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'test111', 'properties.zookeeper.connect' = 'hadoop1:2181', 'properties.bootstrap.servers' = 'hadoop1:9092', 'format' = 'json', 'properties.group.id' = 'iceberg', 'scan.startup.mode' = 'earliest-offset' ); INSERT INTO hive_catalog.test1.sample5 SELECT * FROM default_catalog.default_database.kafka;
查询语句
Batch 模式
sql
SET execution.runtime-mode = batch;
SELECT * FROM sample;
Streaming 模式
sql
SET execution.runtime-mode = streaming;
SET table.dynamic-table-options.enabled = true;
SET sql-client.execution.result-mode = tableau;
-- 从当前快照读取所有记录,然后从该快照读取增量数据
SELECT * FROM sample5 /*+ OPTIONS('streaming'='true', 'monitor-interval'='1s') */;
-- 读取指定快照 id(不包含)后的增量数据
SELECT * FROM sample /*+ OPTIONS('streaming'='true', 'monitor-interval'='1s', 'start-snapshot-id'='3821550127947089987') */;
SQL API 读取 Kafka 数据实时写入 Iceberg 表
-
创建对应的 Iceberg 表
javaStreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); StreamTableEnvironment tblEnv = StreamTableEnvironment.create(env); env.enableCheckpointing(1000); // 1. 创建 Catalog tblEnv.executeSql("CREATE CATALOG hadoop_iceberg WITH (" + "'type'='iceberg'," + "'catalog-type'='hadoop'," + "'warehouse'='hdfs://mycluster/flink_iceberg')"); // 2. 创建 Iceberg 表 flink_iceberg_tbl tblEnv.executeSql("CREATE TABLE hadoop_iceberg.iceberg_db.flink_iceberg_tbl3(id INT, name STRING, age INT, loc STRING) PARTITIONED BY (loc)"); -
编写代码读取 Kafka 数据实时写入 Iceberg
javapublic class ReadKafkaToIceberg { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); StreamTableEnvironment tblEnv = StreamTableEnvironment.create(env); env.enableCheckpointing(1000); // 1. 创建 Catalog tblEnv.executeSql("CREATE CATALOG hadoop_iceberg WITH (" + "'type'='iceberg'," + "'catalog-type'='hadoop'," + "'warehouse'='hdfs://mycluster/flink_iceberg')"); // 2. 创建 Iceberg 表 flink_iceberg_tbl tblEnv.executeSql("CREATE TABLE hadoop_iceberg.iceberg_db.flink_iceberg_tbl3(id INT, name STRING, age INT, loc STRING) PARTITIONED BY (loc)"); // 3. 创建 Kafka Connector,连接消费 Kafka 中数据 tblEnv.executeSql("CREATE TABLE kafka_input_table(" + " id INT," + " name VARCHAR," + " age INT," + " loc VARCHAR" + ") WITH (" + " 'connector' = 'kafka'," + " 'topic' = 'flink-iceberg-topic'," + " 'properties.bootstrap.servers' = 'node1:9092,node2:9092,node3:9092'," + " 'scan.startup.mode' = 'latest-offset'," + " 'properties.group.id' = 'my-group-id'," + " 'format' = 'csv'" + ")"); // 4. 配置 table.dynamic-table-options.enabled Configuration configuration = tblEnv.getConfig().getConfiguration(); configuration.setBoolean("table.dynamic-table-options.enabled", true); // 5. 写入数据到表 flink_iceberg_tbl3 tblEnv.executeSql("INSERT INTO hadoop_iceberg.iceberg_db.flink_iceberg_tbl3 SELECT id, name, age, loc FROM kafka_input_table"); // 6. 查询表数据 TableResult tableResult = tblEnv.executeSql("SELECT * FROM hadoop_iceberg.iceberg_db.flink_iceberg_tbl3 /*+ OPTIONS('streaming'='true', 'monitor-interval'='1s') */"); tableResult.print(); } }
与 Flink 集成的不足
| 支持的特性 | Flink | 备注 |
|---|---|---|
| SQL create catalog | √ | |
| SQL create database | √ | |
| SQL create table | √ | |
| SQL create table like | √ | |
| SQL alter table | √ | 只支持修改表属性,不支持更改列和分区 |
| SQL drop_table | √ | |
| SQL select | √ | 支持流式和批处理模式 |
| SQL insert into | √ | 支持流式和批处理模式 |
| SQL insert overwrite | √ | |
| DataStream read | √ | |
| DataStream append | √ | |
| DataStream overwrite | √ | |
| Metadata tables | 支持 Java API,不支持 Flink SQL | |
| Rewrite files action | √ |
- 不支持创建隐藏分区的 Iceberg 表。
- 不支持创建带有计算列的 Iceberg 表。
- 不支持创建带 watermark 的 Iceberg 表。
- 不支持添加列,删除列,重命名列,更改列。
- Iceberg 目前不支持 Flink SQL 查询表的元数据信息,需要使用 Java API 实现。
大数据学习资料
大数据相关学习资料、大数据项目、湖仓一体、架构师必知必会、数据中台建设方法论...
共有 1400 多份文档资料,另专为星球成员整理了一份比较详细的语雀知识库合计 190 万字和飞书文档资料。(内容太多,仅展示部分内容...),欢迎大家踊跃加入星球,您将获得:
- 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
- 免费答疑解惑、交流技术。
- 面试指导、模拟面试。
- 各类 PDF 文档下载、星球代码下载。
- 提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导。
另外说明加入星球后支持三天无理由退款,不满意无条件随时退。
需要资料请加微信:D1435221412
end
