改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

探索Iceberg与FlinkSQL的深度集成之旅

ckckck2025年1月10日48 浏览

Flink 与 Iceberg 集成指南

环境准备

  1. 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
  2. 上传并解压 Flink 安装包

    bash 复制代码
    tar -zxvf flink-1.16.0-bin-scala_2.12.tgz -C /opt/module/
  3. 配置环境变量

    bash 复制代码
    sudo vim /etc/profile.d/my_env.sh
    export HADOOP_CLASSPATH=hadoop classpath
    source /etc/profile.d/my_env.sh
  4. 拷贝 Iceberg 的 jar 包到 Flink 的 lib 目录

    bash 复制代码
    cp /opt/software/iceberg/iceberg-flink-runtime-1.16-1.1.0.jar /opt/module/flink-1.16.0/lib

启动 Hadoop

启动 SQL-Client

  1. 修改 flink-conf.yaml 配置

    yaml 复制代码
    classloader.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
  2. 启动 Flink

    bash 复制代码
    /opt/module/flink-1.16.0/bin/start-cluster.sh
  3. 启动 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: 内置了 hivehadoop 两种 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

  1. 上传 Hive Connector 到 Flink 的 lib 中

    bash 复制代码
    cp flink-sql-connector-hive-3.1.2_2.12-1.16.0.jar /opt/module/flink-1.16.0/lib/
  2. 启动 Hive Metastore 服务

    bash 复制代码
    hive --service metastore
  3. 创建 Hive Catalog

    sql 复制代码
    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'
    );
    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`;

修改表

  1. 修改表属性

    sql 复制代码
    ALTER TABLE hive_catalog.default.sample SET ('write.format.default'='avro');
  2. 修改表名

    sql 复制代码
    ALTER 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

  1. 建表时指定

    sql 复制代码
    CREATE 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'
    );
  2. 插入时指定

    sql 复制代码
    INSERT INTO tableName /*+ OPTIONS('upsert-enabled'='true') */ ...
  3. 读取 Kafka 流,upsert 插入到 Iceberg 表中

    sql 复制代码
    CREATE 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 表

  1. 创建对应的 Iceberg 表

    java 复制代码
    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)");
  2. 编写代码读取 Kafka 数据实时写入 Iceberg

    java 复制代码
    public 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 备注
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 万字和飞书文档资料。(内容太多,仅展示部分内容...),欢迎大家踊跃加入星球,您将获得:

  1. 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
  2. 免费答疑解惑、交流技术。
  3. 面试指导、模拟面试。
  4. 各类 PDF 文档下载、星球代码下载。
  5. 提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导。

另外说明加入星球后支持三天无理由退款,不满意无条件随时退。

需要资料请加微信:D1435221412

end