Flink SQL 查询引擎连接器

Iceberg

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

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

Apache Iceberg 是面向大规模分析型数据集的开放表格式。Iceberg 通过高性能的表格式,为 Spark、Trino、PrestoDB、Flink、Hive 和 Impala 等计算引擎提供表支持,其使用方式与 SQL 表完全一致。

提示

本文假定你已经掌握了 Iceberg 的基础知识与操作方法。对于本文未涉及的 Iceberg 相关知识,你可以查阅其官方文档。

通过使用 Kyuubi,我们可以直接对 Iceberg 执行 SQL 查询,相比直接使用 Flink 操作 Iceberg 更加便捷、易懂且易于扩展。

Iceberg 集成

要通过 Catalog API 打通 Kyuubi Flink SQL 引擎与 Iceberg 的集成,你需要:

依赖

支持 Iceberg 的 Kyuubi Flink SQL 引擎的 classpath 由以下部分组成:

  1. kyuubi-flink-sql-engine-1.9.1_2.12.jar,随 Kyuubi 发行版部署的引擎 jar 包
  2. 一份 Flink 发行版的拷贝
  3. iceberg-flink-runtime-<flink.version>-<iceberg.version>.jar(例如:iceberg-flink-runtime-1.14-0.14.0.jar),可以在 Maven Central 中找到

为了使引擎的运行时 classpath 能够看到 Iceberg 包,我们可以采用以下任一方式:

  1. 将 Iceberg 包直接放入 $FLINK_HOME/lib
  2. 设置 pipeline.jars=/path/to/iceberg-flink-runtime

警告

请注意不同 Iceberg 与 Flink 版本之间的兼容性,可在 Iceberg 多引擎支持页面上确认。

Iceberg 操作

以 CREATE CATALOG 为例,

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

以 CREATE DATABASE 为例,

CREATE DATABASE iceberg_db;
USE iceberg_db;

以 CREATE TABLE 为例,

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

以 Batch Read 为例,

SET execution.runtime-mode = batch;
SELECT * FROM sample;

以 Streaming Read 为例,

SET execution.runtime-mode = streaming;
SELECT * FROM sample /*+ OPTIONS('streaming'='true', 'monitor-interval'='1s')*/ ;

以 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。

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

评论

登录后参与评论

正在加载评论…