- 数据湖
- 大数据
- 数据存储
【免费下载链接】iceberg
Apache Iceberg
Apache Iceberg 的 Flink 连接器(flink module,仓库内位于flink/v2.3/flink、flink/v2.2/flink、flink/v2.1/flink、flink/v1.20/flink等多个版本目录)提供了完整的 Catalog 接入、批式/流式读取、写入与运维能力。本文以官方配置文档 docs/docs/flink-configuration.md 为骨架,结合连接器源码(如 FlinkCatalogFactory.java、FlinkReadOptions.java、FlinkWriteOptions.java)逐项展开,帮助读者在 Flink SQL 与 DataStream API 两条使用路径上,系统掌握 Catalog 的创建与缓存调优、读/写选项的三种配置渠道与优先级、流式读取起始策略、Range 分布写优化以及table.exec.iceberg.*执行选项。
Catalog 配置
在 Flink SQL 中,Catalog 通过一条CREATE CATALOG语句创建并命名。将<catalog_name>替换为你的 Catalog 名,将<config_key>=<config_value>替换为具体的 Catalog 实现配置项即可:
CREATE CATALOG <catalog_name> WITH ( 'type'='iceberg', `<config_key>`=`<config_value>` );从连接器实现看,type必须为iceberg,它对应 FlinkCatalogFactory.java 中的FACTORY_IDENTIFIER = "iceberg",Flink 依赖该标识定位 Iceberg 的CatalogFactory实现。
通用属性(适用于所有 Catalog 实现)
以下属性对所有 Iceberg Catalog 实现通用,不受具体实现类型限制:
| Property | Required | Values | Description |
|---|---|---|---|
| type | ✔️ | iceberg | 必须为iceberg。 |
| catalog-type | hive、hadoop或rest | 底层 Iceberg Catalog 实现,对应HiveCatalog、HadoopCatalog或RESTCatalog。当通过catalog-impl使用自定义 Catalog 实现(如 AWS Glue、JDBC 或 Nessie Catalog)时必须保持不设置。 | |
| catalog-impl | 自定义 Catalog 实现的完整类名。当catalog-type未设置时必须设置。 | ||
| property-version | 描述属性版本的版本号。该属性用于属性格式变化时的向后兼容。当前属性版本为1。 | ||
| cache-enabled | true或false | 是否启用 Catalog 缓存,默认值为true。 | |
| cache.expiration-interval-ms | Catalog 条目本地缓存的时长,单位为毫秒;负值如-1表示不失效;不允许设置为 0。默认值为-1。 |
catalog-type与catalog-impl互斥:从 FlinkCatalogFactory.java 的createCatalogLoader可以看到,若同时设置了catalog-impl与catalog-type,工厂会抛出IllegalArgumentException("Cannot create catalog ... both catalog-type and catalog-impl are set")。当catalog-impl存在时直接走CatalogLoader.custom(...)自定义加载路径;否则按catalog-type(缺省默认为hive)分派到 Hive / Hadoop / REST 三条内置加载路径。
缓存相关的两个参数同样在工厂中被硬校验:FlinkCatalogFactory.java 通过PropertyUtil.propertyAsBoolean/propertyAsLong解析cache-enabled与cache.expiration-interval-ms,并明确断言cache.expiration-interval-ms不允许为 0。这两个值最终传入FlinkCatalog构造器,决定表、命名空间等 Catalog 元数据在客户端侧的本地缓存行为。
Hive Catalog 属性
| Property | Required | Values | Description |
|---|---|---|---|
| uri | ✔️ | Hive Metastore 的 Thrift URI。 | |
| clients | Hive Metastore 客户端连接池大小,默认值为 2。 | ||
| warehouse | Hive 仓库路径。如果既未通过hive-conf-dir指定包含hive-site.xml的目录,也未在 classpath 中加入正确的hive-site.xml,则用户应显式指定该路径。 | ||
| hive-conf-dir | 包含hive-site.xml配置文件的目录路径,用于提供自定义 Hive 配置。当同时设置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 配置。 |
源码层面,FlinkCatalogFactory.mergeHiveConf 展示了这些配置的真实作用:设置hive-conf-dir时要求该目录下必须存在hive-site.xml(否则抛出IllegalStateException),并将其作为资源加入 HadoopConfiguration;未设置时则尝试从 classpath 加载hive-site.xml,加载不到会在HiveCatalog初始化时抛异常。同理,设置hadoop-conf-dir时要求目录下必须同时存在hdfs-site.xml与core-site.xml。因此「warehouse是否必填」取决于你能否提供正确的hive-site.xml。
Hadoop Catalog 属性
| Property | Required | Values | Description |
|---|---|---|---|
| warehouse | ✔️ | 存放元数据文件和数据文件的 HDFS 目录。 |
REST Catalog 属性
| Property | Required | Values | Description |
|---|---|---|---|
| uri | ✔️ | REST Catalog 的 URL。 | |
| credential | 在 OAuth2 client credentials 流程中用于换取 token 的凭据。 | ||
| token | 用于与服务端交互的 token。 |
REST Catalog 通过CatalogLoader.rest(...)创建,credential与token对应 OAuth2 的两种认证方式,按实际服务端要求二选一或配合使用。除本文表格外,REST Catalog 还支持更多细化参数,可参考 rest-catalog.md。
运行时配置
读选项(Read options)
Flink 读选项在构造 FlinkIcebergSource时传入。以 DataStream API 为例:
IcebergSource.forRowData() .tableLoader(TableLoader.fromCatalog(...)) .assignerFactory(new SimpleSplitAssignerFactory()) .streaming(true) .streamingStartingStrategy(StreamingStartingStrategy.INCREMENTAL_FROM_SNAPSHOT_ID) .startSnapshotId(3821550127947089987L) .monitorInterval(Duration.ofMillis(10L)) // 或 .set("monitor-interval", "10s") / set(FlinkReadOptions.MONITOR_INTERVAL, "10s") .build()在 Flink SQL 中,读选项可以通过 SQL Hint 传入:
SELECT * FROM tableName /*+ OPTIONS('monitor-interval'='10s') */ ...选项还可以通过 Flink 配置(Flink configuration)传入,并作用于当前会话。注意:并非所有选项都支持这种方式:
env.getConfig() .getConfiguration() .set(FlinkReadOptions.SPLIT_FILE_OPEN_COST_OPTION, 1000L); ...优先级:Read option(读选项)优先级最高,其次是Flink configuration(Flink 配置),最后才是Table property(表属性)。
这一优先级在 FlinkReadConf.java 与 FlinkConfParser.java 中落地:例如splitSize()的解析链是「读选项 → Flink 配置(connector.iceberg.split-size)→ 表属性(read.split.target-size)→ 默认值(128 MB)」。此外,FlinkReadOptions.java 中 split 相关选项特意不声明默认值,正是为了避免默认值掩盖FlinkReadConf中对表属性的回退逻辑。
完整的读选项如下:
| Read option | Flink configuration | Table property | Default | Description |
|---|---|---|---|---|
| snapshot-id | N/A | N/A | null | 批模式下做时间旅行(time travel),从指定 snapshot-id 读取数据。 |
| case-sensitive | connector.iceberg.case-sensitive | N/A | false | 为 true 时按大小写敏感方式匹配列名。 |
| as-of-timestamp | N/A | N/A | null | 批模式下做时间旅行,从给定时间(毫秒)最近的一个 snapshot 读取数据。 |
| starting-strategy | connector.iceberg.starting-strategy | N/A | INCREMENTAL_FROM_LATEST_SNAPSHOT | 流式执行的起始策略。TABLE_SCAN_THEN_INCREMENTAL:先做一次常规表扫描再切换为增量模式,增量模式从当前 snapshot 之后(exclusive)开始。INCREMENTAL_FROM_LATEST_SNAPSHOT:从最新 snapshot(含,inclusive)开始增量模式;若为空表,则发现所有未来的 append snapshot。INCREMENTAL_FROM_LATEST_SNAPSHOT_EXCLUSIVE:从最新 snapshot 之后(exclusive)开始增量模式;若为空表,则发现所有未来的 append snapshot。INCREMENTAL_FROM_EARLIEST_SNAPSHOT:从最早 snapshot(含)开始增量模式;若为空表,则发现所有未来的 append snapshot。INCREMENTAL_FROM_SNAPSHOT_ID:从指定 id 的 snapshot(含)开始增量模式。INCREMENTAL_FROM_SNAPSHOT_TIMESTAMP:从指定时间戳的 snapshot(含)开始增量模式;若时间戳位于两个 snapshot 之间,则从该时间戳之后的 snapshot 开始。该策略仅适用于 FIP-27 Source。 |
| start-snapshot-timestamp | N/A | N/A | null | 从给定时间(毫秒)最近的一个 snapshot 开始读取数据。 |
| start-snapshot-id | N/A | N/A | null | 从指定 snapshot-id 开始读取数据。 |
| end-snapshot-id | N/A | N/A | The latest snapshot id | 指定结束 snapshot。 |
| branch | N/A | N/A | main | 批模式下指定要读取的 branch。 |
| tag | N/A | N/A | null | 批模式下指定要读取的 tag。 |
| start-tag | N/A | N/A | null | 增量读取时指定起始 tag。 |
| end-tag | N/A | N/A | null | 增量读取时指定结束 tag。 |
| split-size | connector.iceberg.split-size | read.split.target-size | 128 MB | 合并输入 split 时的目标大小。 |
| split-lookback | connector.iceberg.split-lookback | read.split.planning-lookback | 10 | 合并输入 split 时考虑的 bin 数量。 |
| split-file-open-cost | connector.iceberg.split-file-open-cost | read.split.open-file-cost | 4MB | 打开文件的预估开销,作为合并 split 时的最小权重。 |
| streaming | connector.iceberg.streaming | N/A | false | 设置当前任务运行在流式还是批式模式。 |
| monitor-interval | connector.iceberg.monitor-interval | N/A | 60s | 发现新 snapshot split 的监控间隔。仅适用于流式读取。 |
| include-column-stats | connector.iceberg.include-column-stats | N/A | false | 创建新扫描时随每个数据文件加载列统计信息。列统计包括:value count、null value count、lower bounds 和 upper bounds。 |
| max-planning-snapshot-count | connector.iceberg.max-planning-snapshot-count | N/A | Integer.MAX_VALUE | 每次 split 枚举最多限制的 snapshot 数量。仅适用于流式读取。 |
| limit | connector.iceberg.limit | N/A | -1 | 限制输出的行数。 |
| max-allowed-planning-failures | connector.iceberg.max-allowed-planning-failures | N/A | 3 | 扫描规划失败前允许的最大连续失败次数。设为 -1 表示扫描规划失败也永不使作业失败。 |
| watermark-column | connector.iceberg.watermark-column | N/A | null | 指定用于生成 watermark 的列。该选项存在时,splitAssignerFactory会被覆盖为OrderedSplitAssignerFactory。 |
| watermark-column-time-unit | connector.iceberg.watermark-column-time-unit | N/A | TimeUnit.MICROSECONDS | 指定生成 watermark 使用的时间单位。可选值:DAYS、HOURS、MINUTES、SECONDS、MILLISECONDS、MICROSECONDS、NANOSECONDS。 |
从源码确认的几点实现细节:
- 流式起始策略定义在 StreamingStartingStrategy.java 中,共 6 个枚举值,默认值为
INCREMENTAL_FROM_LATEST_SNAPSHOT,对应 FlinkReadOptions.java 中STARTING_STRATEGY_OPTION的默认值。 watermark-column-time-unit在 FlinkReadOptions.java 中直接映射为 Java 的TimeUnit枚举类型。split-size、split-lookback、split-file-open-cost三个选项除了读选项与 Flink 配置渠道外,还支持表属性回退(对应read.split.target-size、read.split.planning-lookback、read.split.open-file-cost),这是表格中唯一带「Table property」一列的读选项。- 时间旅行与 branch/tag 读取均只在批模式下生效;增量读取的
start-tag/end-tag、start-snapshot-id/end-snapshot-id等则用于流式增量场景。关于 branch、tag 的更完整语义可参考 branching.md。
写选项(Write options)
Flink 写选项在构造FlinkSink时传入,例如:
FlinkSink.Builder builder = FlinkSink.forRow(dataStream, SimpleDataUtil.FLINK_SCHEMA) .table(table) .tableLoader(tableLoader) .set("write-format", "orc") .set(FlinkWriteOptions.OVERWRITE_MODE, "true");Flink SQL 中通过 SQL Hint 传入写选项:
INSERT INTO tableName /*+ OPTIONS('upsert-enabled'='true') */ ...写选项清单如下:
| Flink option | Default | Description |
|---|---|---|
| write-format | Table write.format.default | 本次写入使用的文件格式:parquet、avro 或 orc |
| target-file-size-bytes | As per table property | 覆盖该表的 write.target-file-size-bytes |
| upsert-enabled | Table write.upsert.enabled | 覆盖该表的 write.upsert.enabled |
| overwrite-enabled | false | 覆盖表数据;配置使用 UPSERT 数据流时不应启用覆盖模式。 |
| distribution-mode | Table write.distribution-mode | 覆盖该表的 write.distribution-mode。RANGE 分布目前处于实验状态。 |
| range-distribution-statistics-type | Auto | Range 分布的数据统计收集类型:Map、Sketch、Auto。详见下文「Range distribution statistics type」。 |
| range-distribution-sort-key-base-weight | 0.0 (double) | 每个排序键相对每个 writer task 目标流量权重的基准权重。详见下文「Range distribution sort key base weight」。 |
| compression-codec | Table write.(fileformat).compression-codec | 覆盖该表本次写入的压缩编码 |
| compression-level | Table write.(fileformat).compression-level | 覆盖该表本次写入的 Parquet 与 Avro 压缩级别 |
| compression-strategy | Table write.orc.compression-strategy | 覆盖该表本次写入的 ORC 压缩策略 |
| write-parallelism | Upstream operator parallelism | 覆盖写入并行度 |
| branch | main | 写入的 branch |
| uid-suffix | As per table property | 覆盖该表底层 IcebergSink 使用的 uid suffix |
| shred-variants | Table write.parquet.shred-variants | 覆盖该表本次写入的 variant shred 配置 |
| variant-inference-buffer-size | Table write.parquet.variant-inference-buffer-size | 覆盖该表本次写入的 variant 推断缓冲区大小 |
| flink-maintenance.rewrite.enabled | false | 提交成功后运行数据文件压缩(compaction)。仅由IcebergSink使用,见 flink-maintenance.md 的 IcebergSink with post-commit integration 一节。 |
| flink-maintenance.expire-snapshots.enabled | false | 提交成功后清理过期 snapshot。仅由IcebergSink使用,同上。 |
| flink-maintenance.delete-orphan-files.enabled | false | 提交成功后删除孤儿文件。仅由IcebergSink使用,同上。 |
| flink-maintenance.convert-equality-deletes.enabled | false | 提交成功后把 equality delete 转换为 deletion vectors。仅由IcebergSink使用,同上。 |
源码层面的补充说明:
- FlinkWriteOptions.java 中
overwrite-enabled有显式默认值false,而多数覆盖表属性的选项(如write-format、target-file-size-bytes、compression-codec等)采用noDefaultValue(),即未设置时回退到表属性。 branch的默认值来自SnapshotRef.MAIN_BRANCH(即main),见 FlinkWriteOptions.java。- 后提交维护(post-commit maintenance)的四个
flink-maintenance.*.enabled开关分别对应 FlinkWriteOptions.java 中RewriteDataFilesConfig.PREFIX、ExpireSnapshotsConfig.PREFIX、DeleteOrphanFilesConfig.PREFIX、ConvertEqualityDeletesConfig.PREFIX派生的配置键,且rewrite.enabled保留了compaction-enabled作为弃用键(deprecated key)。详细用法见 flink-maintenance.md。
Range 分布统计类型(Range distribution statistics type)
该配置值是枚举类型:Map、Sketch、Auto。
- Map:为每个键收集精确的采样计数。适用于低基数场景(如成百上千个键)。
- Sketch:通过蓄水池采样(reservoir sampling)构造均匀随机采样。适合高基数场景(如百万级),因为内存占用保持很低。
- Auto:以 Map 统计开始;一旦检测到基数超过阈值(当前为 10,000),自动切换为 Sketch。
从 FlinkWriteOptions.java 可见range-distribution-statistics-type的默认值为StatisticsType.Auto。该选项作用于flink.sink.shuffle包下的 Range 分布实现(实验特性),用于为每个排序键的流量统计提供数据基础。
Range 分布排序键基准权重(Range distribution sort key base weight)
range-distribution-sort-key-base-weight:0.0(double 类型)。
若排序顺序包含分区列,每个排序键会映射到一个分区和一个数据文件。该相对权重可以避免为低流量的排序键产生过多小文件。它是一个 double 值,定义每个排序键的最小权重。2.0表示每个键具有每个 writer task 目标流量权重2%的基准权重。
示例:假设 sink Iceberg 表按事件时间每日分区。数据流包含从现在到 180 天前的事件。按事件时间划分时,不同日期间的流量权重分布通常呈长尾形态——当天流量最大,越老的日期(长尾)流量越少。假设 writer 并行度为10,180 天的总权重为10,000,则每个 writer task 的目标流量权重为1,000。假设最老的 150 天权重总和为1,000,正常情况下 Range 分区器会把最老的 150 天全部放在一个 writer task 上,该 task 将写出 150 个小文件(每天一个)。保持 150 个打开的文件可能消耗大量内存;checkpoint 时 flush 并上传 150 个文件(无论多小)也可能很慢。若将该配置设为2.0,意味着每个排序键都有目标权重1,000的2%基准权重,这样无论数据多小,都能避免在单个 writer task 上放置超过50个数据文件(每天一个)。
此配置仅适用于低基数场景的StatisticsType.Map。对于StatisticsType.Sketch高基数排序列,通常不会用作分区列,否则写入时可能产生过多分区和小文件——Sketch Range 分区器只是把高基数键拆分成有序区间。
执行选项(Execution options)
Iceberg 还提供一组table.exec.iceberg.*选项,它们从 Flink 配置中读取,而非作为每次作业的读/写选项。SQL 中通过会话级设置:
SET 'table.exec.iceberg.infer-source-parallelism' = 'false';使用 DataStream API 时,可以在传给 source 或 sink builder 的Configuration中设置。例如 FlinkConfigOptions.java 的类注释所示:
Configuration configuration = new Configuration(); configuration.setBoolean(FlinkConfigOptions.TABLE_EXEC_ICEBERG_INFER_SOURCE_PARALLELISM, true); FlinkSource.forRowData() .flinkConf(configuration) ...执行选项清单如下:
| Flink configuration | Default | Description |
|---|---|---|
| table.exec.iceberg.infer-source-parallelism | true | 为 true 时,批读取的 source 并行度由 scan split 数量推断得出,上限为table.exec.iceberg.infer-source-parallelism.max,且受查询 limit(若设置)约束;为 false 时,source 并行度取自 Flink 配置。流式读取从不推断并行度。 |
| table.exec.iceberg.infer-source-parallelism.max | 100 | source 算子推断并行度的最大值。 |
| table.exec.iceberg.fetch-batch-record-count | 2048 | Iceberg source reader 每次 fetch 批次的目标准入记录数。 |
| table.exec.iceberg.worker-pool-size | max(2, available cpu) | 用于规划或扫描 manifest 的 worker 池大小。默认为共享 Iceberg worker 池大小,受iceberg.worker.num-threads系统属性控制。 |
| table.exec.iceberg.use-v2-sink | false | 使用基于 SinkV2 的IcebergSink实现,见 flink-writes.md 的 Sink V2 based implementation 一节。 |
这些选项的定义集中在 FlinkConfigOptions.java 中,例如:
table.exec.iceberg.infer-source-parallelism(默认true)与.max(默认100)对应 FlinkConfigOptions.java,批模式下根据 scan split 数量推断并行度,流式读取不推断。table.exec.iceberg.fetch-batch-record-count默认2048(FlinkConfigOptions.java),控制 source reader 每个 fetch 批次的记录数目标。table.exec.iceberg.worker-pool-size默认取ThreadPools.WORKER_THREAD_POOL_SIZE(FlinkConfigOptions.java),即共享 Iceberg worker 池大小,可通过iceberg.worker.num-threads系统属性调整。table.exec.iceberg.use-v2-sink默认false(FlinkConfigOptions.java),开启后切换到 SinkV2 API 实现的IcebergSink。
此外该文件还定义了三个文档主表未列出但同样以table.exec.iceberg.为前缀的选项:table.exec.iceberg.expose-split-locality-info(暴露 split 主机信息以使用 Flink 的 locality 感知 split assigner)、table.exec.iceberg.use-flip27-source(默认true,使用 FLIP-27 的 Iceberg source 实现)以及table.exec.iceberg.split-assigner-type(默认SIMPLE,决定 split 如何分配给 reader)。这些选项与文档表格中的五个执行选项共同构成连接器的完整执行层调优面。
配置优先级总结
综合全文,Iceberg Flink 连接器的配置遵循以下层次:
- 读选项 / 写选项(Read/Write option):DataStream API 的 builder
.set(...)或 SQL 的/*+ OPTIONS(...) */Hint,优先级最高。 - Flink 配置(Flink configuration):会话级
SET或env.getConfig().getConfiguration()中的connector.iceberg.*/table.exec.iceberg.*配置。其中connector.iceberg.*面向读/写行为(并非所有读选项都支持该渠道),table.exec.iceberg.*面向执行行为。 - 表属性(Table property):仅部分读选项(如 split 系列)与多数写选项(如
write-format、target-file-size-bytes、distribution-mode等)支持回退到表属性,作为未显式指定时的兜底默认。
理解这一优先级模型,可以在集群级默认、会话级覆盖、作业级定制三个粒度上灵活管理 Iceberg + Flink 的读写行为:集群默认值写在表属性或 Flink 配置中,个别作业需要差异化时再用读/写选项按作业覆盖,而无需改动表定义或全局配置。
- 数据湖
- 大数据
- 数据存储
【免费下载链接】iceberg
Apache Iceberg
相关推荐
Flink DataStream 执行配置详解:ExecutionConfig 全选项与源码级原理解析
Flink DataStream 执行配置详解:ExecutionConfig 全选项与源码级原理解析 StreamExecutionEnvironment 内
大数据流处理批处理数据工程Apache Spark SQL XML 数据源完全指南:读取、写入与选项详解
Apache Spark SQL XML 数据源完全指南:读取、写入与选项详解 导读 本文围绕 Apache Spark SQL 内置的 XML 数据源( sp
大数据数据分析批处理流处理机器学习图计算Robot Framework 命令行选项完全指南:robot 执行与 rebot 后处理选项及环境变量详解
Robot Framework 命令行选项完全指南:robot 执行与 rebot 后处理选项及环境变量详解 导读 本篇指南系统梳理 Robot Framewo
测试RPA接口测试
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考