Apache Flink

Flink DDL

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

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

Flink DDL

DDL 命令

CREATE Catalog

Hive catalog

这会创建一个名为 hive_catalog 的 Iceberg catalog,可通过 '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'
);

使用 Hive catalog 时,可以设置以下属性:

  • uri:Hive metastore 的 thrift URI。(必填)
  • clients:Hive metastore 客户端连接池大小,默认值为 2。(可选)
  • warehouse:Hive warehouse 的位置。如果既没有将 hive-conf-dir 设置为包含 hive-site.xml 配置文件的目录,也没有将正确的 hive-site.xml 添加到 classpath 中,则用户应指定此路径。
  • hive-conf-dir:包含 hive-site.xml 配置文件的目录路径,该文件用于提供自定义的 Hive 配置值。创建 iceberg catalog 时,如果同时设置了 hive-conf-dir 和 warehouse,则 <hive-conf-dir>/hive-site.xml(或 classpath 中的 hive 配置文件)中的 hive.metastore.warehouse.dir 的值将被 warehouse 的值覆盖。
  • hadoop-conf-dir:包含 core-site.xml 和 hdfs-site.xml 配置文件的目录路径,这些文件用于提供自定义的 Hadoop 配置值。

Hive Catalog 的限制

Hive Metastore(HMS)会通过按位置比较列类型来验证 schema 变更(hive.metastore.disallow.incompatible.col.type.changes,默认为 true)。使用 Hive catalog 时,改变列位置的 schema 演进操作——例如删除非最后一列或重新排列列顺序——可能会失败,无论执行该操作的是哪个引擎(Spark、Flink Java API 等)。

要规避此问题,可将 hive.metastore.disallow.incompatible.col.type.changes=false 以禁用 HMS 的 schema 兼容性检查:

  • 远程 HMS: 在 HMS 服务器的 hive-site.xml 中设置此属性。
  • 嵌入式 HMS: 在 Hive catalog 配置中添加等效属性。

权衡: 禁用此检查后,由于 Hive Metastore 中的 schema 不匹配,Hive 引擎可能无法再正确读取该表。支持 Iceberg 的引擎(Spark、Flink、Trino 等)将继续正常工作,因为它们从 Iceberg 元数据而非 Hive Metastore 中读取 schema。

Hadoop catalog

Iceberg 还支持 HDFS 中基于目录的 catalog,可通过 'catalog-type'='hadoop' 进行配置:

CREATE CATALOG hadoop_catalog WITH (
  'type'='iceberg',
  'catalog-type'='hadoop',
  'warehouse'='hdfs://nn:8020/warehouse/path',
  'property-version'='1'
);

如果使用 Hadoop catalog,可以设置以下属性:

  • warehouse:存储元数据文件和数据文件的 HDFS 目录。(必填)

执行 SQL 命令 USE CATALOG hadoop_catalog 来设置当前 catalog。

REST catalog

这会创建一个名为 rest_catalog 的 iceberg catalog,可以通过 'catalog-type'='rest' 进行配置,用于从 REST catalog 加载表:

CREATE CATALOG rest_catalog WITH (
  'type'='iceberg',
  'catalog-type'='rest',
  'uri'='https://localhost/'
);

如果使用 REST catalog,可以设置以下属性:

  • uri:REST Catalog 的 URL(必填)
  • credential:在 OAuth2 客户端凭据流程中用于换取令牌的凭据(可选)
  • token:用于与服务端交互的令牌(可选)

自定义 catalog

Flink 还支持通过指定 catalog-impl 属性来加载自定义的 Iceberg Catalog 实现:

CREATE CATALOG my_catalog WITH (
  'type'='iceberg',
  'catalog-impl'='com.my.custom.CatalogImpl',
  'my-additional-catalog-config'='my-value'
);

通过 YAML 配置创建

在启动 SQL 客户端之前,可以在 sql-client-defaults.yaml 中注册目录(Catalog)。

catalogs:
  - name: my_catalog
    type: iceberg
    catalog-type: hadoop
    warehouse: hdfs://nn:8020/warehouse/path

通过 SQL 文件创建

Flink SQL Client 支持 -i 启动选项,用于在启动 SQL Client 时执行初始化 SQL 文件以设置环境。

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

USE CATALOG hive_catalog;

使用 -i <init.sql> 选项来初始化 SQL Client 会话:

/path/to/bin/sql-client.sh -i /path/to/init.sql

CREATE DATABASE

默认情况下,Iceberg 在 Flink 中会使用 default 数据库。可以使用下面的示例创建一个独立的数据库,以避免将表创建在 default 数据库下:

CREATE DATABASE iceberg_db;
USE iceberg_db;

CREATE TABLE 创建表

CREATE TABLE `hive_catalog`.`default`.`sample` (
    id BIGINT COMMENT 'unique id',
    data STRING NOT NULL
) WITH ('format-version'='2');

建表语句支持常用的 Flink create 子句,包括:

  • PARTITION BY (column1, column2, ...) 用于配置分区,Flink 目前尚不支持隐藏分区。
  • COMMENT 'table document' 用于设置表说明。
  • WITH ('key'='value', ...) 用于设置表配置,这些配置会保存在 Iceberg 表属性中。

要指定表的位置,请使用 WITH ('location'='fully-qualified-uri'):

CREATE TABLE `hive_catalog`.`default`.`sample` (
    id BIGINT COMMENT 'unique id',
    data STRING NOT NULL
) WITH (
    'format-version'='2',
    'location'='hdfs//nn:8020/custom-path'
);

目前,它不支持计算列、水印定义等。

PRIMARY KEY

可以为单个列或一组列声明主键约束,主键值必须唯一且不能为空。UPSERT 模式要求必须定义主键。

CREATE TABLE `hive_catalog`.`default`.`sample` (
    id BIGINT COMMENT 'unique id',
    data STRING NOT NULL,
    PRIMARY KEY(`id`) NOT ENFORCED
) WITH ('format-version'='2');

PARTITIONED BY

要创建分区表,请使用 PARTITIONED BY:

CREATE TABLE `hive_catalog`.`default`.`sample` (
    id BIGINT COMMENT 'unique id',
    data STRING NOT NULL
)
PARTITIONED BY (data)
WITH ('format-version'='2');

Iceberg 支持隐藏分区,但 Flink 不支持按列上的函数进行分区,因此无法在 Flink DDL 中支持隐藏分区。

CREATE TABLE LIKE

若要创建一个与另一张表具有相同 schema、分区方式和表属性的表,可以使用 CREATE TABLE LIKE。

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

CREATE TABLE  `hive_catalog`.`default`.`sample_like` LIKE `hive_catalog`.`default`.`sample`;

更多详情,请参阅 Flink CREATE TABLE 文档。

ALTER TABLE

Iceberg 仅支持修改表属性:

ALTER TABLE `hive_catalog`.`default`.`sample` SET ('write.format.default'='avro');

ALTER TABLE .. RENAME TO

ALTER TABLE `hive_catalog`.`default`.`sample` RENAME TO `hive_catalog`.`default`.`new_sample`;

DROP TABLE

要删除表,请执行:

DROP TABLE `hive_catalog`.`default`.`sample`;

评论

登录后参与评论

正在加载评论…