Apache Flink

Flink 连接器

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

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

Flink Connector

Apache Flink 支持直接创建 Iceberg 表,而无需在 Flink SQL 中显式创建 Flink catalog。这意味着我们可以在 Flink SQL 中通过指定表选项 'connector'='iceberg' 来创建 Iceberg 表,其用法与 Flink 官方文档中的用法类似。

在 Flink 中,SQL 语句 CREATE TABLE test (..) WITH ('connector'='iceberg', ...) 会在当前 Flink catalog 中创建一个 Flink 表(默认使用 GenericInMemoryCatalog),该表只是映射到底层的 Iceberg 表,而不是在当前 Flink catalog 中直接维护 Iceberg 表。

要通过 SQL 语法 CREATE TABLE test (..) WITH ('connector'='iceberg', ...) 在 Flink SQL 中创建表,Flink Iceberg connector 允许通过表属性来设置 catalog 属性。各属性的有效值详见 Flink Configuration 页面。

Hive catalog 中管理的表

在执行以下 SQL 之前,请确保已按照快速开始文档正确配置了 Flink SQL client。

以下 SQL 会在当前 Flink catalog 中创建一个 Flink 表,该表映射到 Iceberg catalog 中管理的 Iceberg 表 default_database.flink_table。

CREATE TABLE flink_table (
    id   BIGINT,
    data STRING
) WITH (
    'connector'='iceberg',
    'catalog-name'='hive_prod',
    'uri'='thrift://localhost:9083',
    'warehouse'='hdfs://nn:8020/path/to/warehouse'
);

如果你想创建一个 Flink 表,用于映射到 Hive catalog 中管理的另一个 Iceberg 表(例如 Hive 中的 hive_db.hive_iceberg_table),那么可以按如下方式创建 Flink 表:

CREATE TABLE flink_table (
    id   BIGINT,
    data STRING
) WITH (
    'connector'='iceberg',
    'catalog-name'='hive_prod',
    'catalog-database'='hive_db',
    'catalog-table'='hive_iceberg_table',
    'uri'='thrift://localhost:9083',
    'warehouse'='hdfs://nn:8020/path/to/warehouse'
);

信息

在向 Flink 表写入记录时,如果底层的 catalog 数据库(上例中的 hive_db)不存在,将会被自动创建。

hadoop catalog 中托管的表

以下 SQL 将在当前 Flink catalog 中创建一个 Flink 表,该表映射到 hadoop catalog 中托管的 iceberg 表 default_database.flink_table。

CREATE TABLE flink_table (
    id   BIGINT,
    data STRING
) WITH (
    'connector'='iceberg',
    'catalog-name'='hadoop_prod',
    'catalog-type'='hadoop',
    'warehouse'='hdfs://nn:8020/path/to/warehouse'
);

REST Catalog 中托管的表

以下 SQL 将在当前 Flink catalog 中创建一个 Flink 表,该表映射到 REST Catalog 中托管的 Iceberg 表 default_database.flink_table。

CREATE TABLE flink_table (
    id   BIGINT,
    data STRING
) WITH (
    'connector'='iceberg',
    'catalog-name'='rest_prod',
    'catalog-type'='rest',
    'uri'='https://localhost/'
    'credential'='xxxx' -- Optional
    'token'='xxxx' -- Optional
    'scope'='xxxx' -- Optional
     ...
);

在自定义 catalog 中管理的表

以下 SQL 会在当前 Flink catalog 中创建一张 Flink 表,该表映射到在类型为 com.my.custom.CatalogImpl 的自定义 catalog 中管理的 iceberg 表 default_database.flink_table。

CREATE TABLE flink_table (
    id   BIGINT,
    data STRING
) WITH (
    'connector'='iceberg',
    'catalog-name'='custom_prod',
    'catalog-impl'='com.my.custom.CatalogImpl',
     -- More table properties for the customized catalog
    'my-additional-catalog-config'='my-value',
     ...
);

请查看「Integrations」选项卡下所有自定义目录的相关章节。

完整示例。

以 Hive 目录为例:

CREATE TABLE flink_table (
    id   BIGINT,
    data STRING
) WITH (
    'connector'='iceberg',
    'catalog-name'='hive_prod',
    'uri'='thrift://localhost:9083',
    'warehouse'='file:///path/to/warehouse'
);

INSERT INTO flink_table VALUES (1, 'AAA'), (2, 'BBB'), (3, 'CCC');

SET execution.result-mode=tableau;
SELECT * FROM flink_table;

+----+------+
| id | data |
+----+------+
|  1 |  AAA |
|  2 |  BBB |
|  3 |  CCC |
+----+------+
3 rows in set

更多详情,请参阅 Iceberg 的 Flink 文档。

评论

登录后参与评论

正在加载评论…