ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Iceberg Flink 配置完全指南:Catalog、读选项、写选项与执行选项详解

Apache Iceberg Flink 配置完全指南:Catalog、读选项、写选项与执行选项详解 数据湖大数据数据存储【免费下载链接】icebergApache Iceberg项目地址https://gitcode.com/gh_mirrors/icebe/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_keyconfig_value替换为具体的 Catalog 实现配置项即可CREATE CATALOG catalog_name WITH ( typeiceberg, config_keyconfig_value );从连接器实现看type必须为iceberg它对应 FlinkCatalogFactory.java 中的FACTORY_IDENTIFIER icebergFlink 依赖该标识定位 Iceberg 的CatalogFactory实现。通用属性适用于所有 Catalog 实现以下属性对所有 Iceberg Catalog 实现通用不受具体实现类型限制PropertyRequiredValuesDescriptiontype✔️iceberg必须为iceberg。catalog-typehive、hadoop或rest底层 Iceberg Catalog 实现对应HiveCatalog、HadoopCatalog或RESTCatalog。当通过catalog-impl使用自定义 Catalog 实现如 AWS Glue、JDBC 或 Nessie Catalog时必须保持不设置。catalog-impl自定义 Catalog 实现的完整类名。当catalog-type未设置时必须设置。property-version描述属性版本的版本号。该属性用于属性格式变化时的向后兼容。当前属性版本为1。cache-enabledtrue或false是否启用 Catalog 缓存默认值为true。cache.expiration-interval-msCatalog 条目本地缓存的时长单位为毫秒负值如-1表示不失效不允许设置为 0。默认值为-1。catalog-type与catalog-impl互斥从 FlinkCatalogFactory.java 的createCatalogLoader可以看到若同时设置了catalog-impl与catalog-type工厂会抛出IllegalArgumentExceptionCannot 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 属性PropertyRequiredValuesDescriptionuri✔️Hive Metastore 的 Thrift URI。clientsHive Metastore 客户端连接池大小默认值为 2。warehouseHive 仓库路径。如果既未通过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 属性PropertyRequiredValuesDescriptionwarehouse✔️存放元数据文件和数据文件的 HDFS 目录。REST Catalog 属性PropertyRequiredValuesDescriptionuri✔️REST Catalog 的 URL。credential在 OAuth2 client credentials 流程中用于换取 token 的凭据。token用于与服务端交互的 token。REST Catalog 通过CatalogLoader.rest(...)创建credential与token对应 OAuth2 的两种认证方式按实际服务端要求二选一或配合使用。除本文表格外REST Catalog 还支持更多细化参数可参考 rest-catalog.md。运行时配置读选项Read optionsFlink 读选项在构造 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-interval10s) */ ...选项还可以通过 Flink 配置Flink configuration传入并作用于当前会话。注意并非所有选项都支持这种方式env.getConfig() .getConfiguration() .set(FlinkReadOptions.SPLIT_FILE_OPEN_COST_OPTION, 1000L); ...优先级Read option读选项优先级最高其次是Flink configurationFlink 配置最后才是Table property表属性。这一优先级在 FlinkReadConf.java 与 FlinkConfParser.java 中落地例如splitSize()的解析链是「读选项 → Flink 配置connector.iceberg.split-size→ 表属性read.split.target-size→ 默认值128 MB」。此外FlinkReadOptions.java 中 split 相关选项特意不声明默认值正是为了避免默认值掩盖FlinkReadConf中对表属性的回退逻辑。完整的读选项如下Read optionFlink configurationTable propertyDefaultDescriptionsnapshot-idN/AN/Anull批模式下做时间旅行time travel从指定 snapshot-id 读取数据。case-sensitiveconnector.iceberg.case-sensitiveN/Afalse为 true 时按大小写敏感方式匹配列名。as-of-timestampN/AN/Anull批模式下做时间旅行从给定时间毫秒最近的一个 snapshot 读取数据。starting-strategyconnector.iceberg.starting-strategyN/AINCREMENTAL_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-timestampN/AN/Anull从给定时间毫秒最近的一个 snapshot 开始读取数据。start-snapshot-idN/AN/Anull从指定 snapshot-id 开始读取数据。end-snapshot-idN/AN/AThe latest snapshot id指定结束 snapshot。branchN/AN/Amain批模式下指定要读取的 branch。tagN/AN/Anull批模式下指定要读取的 tag。start-tagN/AN/Anull增量读取时指定起始 tag。end-tagN/AN/Anull增量读取时指定结束 tag。split-sizeconnector.iceberg.split-sizeread.split.target-size128 MB合并输入 split 时的目标大小。split-lookbackconnector.iceberg.split-lookbackread.split.planning-lookback10合并输入 split 时考虑的 bin 数量。split-file-open-costconnector.iceberg.split-file-open-costread.split.open-file-cost4MB打开文件的预估开销作为合并 split 时的最小权重。streamingconnector.iceberg.streamingN/Afalse设置当前任务运行在流式还是批式模式。monitor-intervalconnector.iceberg.monitor-intervalN/A60s发现新 snapshot split 的监控间隔。仅适用于流式读取。include-column-statsconnector.iceberg.include-column-statsN/Afalse创建新扫描时随每个数据文件加载列统计信息。列统计包括value count、null value count、lower bounds 和 upper bounds。max-planning-snapshot-countconnector.iceberg.max-planning-snapshot-countN/AInteger.MAX_VALUE每次 split 枚举最多限制的 snapshot 数量。仅适用于流式读取。limitconnector.iceberg.limitN/A-1限制输出的行数。max-allowed-planning-failuresconnector.iceberg.max-allowed-planning-failuresN/A3扫描规划失败前允许的最大连续失败次数。设为 -1 表示扫描规划失败也永不使作业失败。watermark-columnconnector.iceberg.watermark-columnN/Anull指定用于生成 watermark 的列。该选项存在时splitAssignerFactory会被覆盖为OrderedSplitAssignerFactory。watermark-column-time-unitconnector.iceberg.watermark-column-time-unitN/ATimeUnit.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 optionsFlink 写选项在构造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-enabledtrue) */ ...写选项清单如下Flink optionDefaultDescriptionwrite-formatTable write.format.default本次写入使用的文件格式parquet、avro 或 orctarget-file-size-bytesAs per table property覆盖该表的 write.target-file-size-bytesupsert-enabledTable write.upsert.enabled覆盖该表的 write.upsert.enabledoverwrite-enabledfalse覆盖表数据配置使用 UPSERT 数据流时不应启用覆盖模式。distribution-modeTable write.distribution-mode覆盖该表的 write.distribution-mode。RANGE 分布目前处于实验状态。range-distribution-statistics-typeAutoRange 分布的数据统计收集类型Map、Sketch、Auto。详见下文「Range distribution statistics type」。range-distribution-sort-key-base-weight0.0 (double)每个排序键相对每个 writer task 目标流量权重的基准权重。详见下文「Range distribution sort key base weight」。compression-codecTable write.(fileformat).compression-codec覆盖该表本次写入的压缩编码compression-levelTable write.(fileformat).compression-level覆盖该表本次写入的 Parquet 与 Avro 压缩级别compression-strategyTable write.orc.compression-strategy覆盖该表本次写入的 ORC 压缩策略write-parallelismUpstream operator parallelism覆盖写入并行度branchmain写入的 branchuid-suffixAs per table property覆盖该表底层 IcebergSink 使用的 uid suffixshred-variantsTable write.parquet.shred-variants覆盖该表本次写入的 variant shred 配置variant-inference-buffer-sizeTable write.parquet.variant-inference-buffer-size覆盖该表本次写入的 variant 推断缓冲区大小flink-maintenance.rewrite.enabledfalse提交成功后运行数据文件压缩compaction。仅由IcebergSink使用见 flink-maintenance.md 的 IcebergSink with post-commit integration 一节。flink-maintenance.expire-snapshots.enabledfalse提交成功后清理过期 snapshot。仅由IcebergSink使用同上。flink-maintenance.delete-orphan-files.enabledfalse提交成功后删除孤儿文件。仅由IcebergSink使用同上。flink-maintenance.convert-equality-deletes.enabledfalse提交成功后把 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 weightrange-distribution-sort-key-base-weight0.0double 类型。若排序顺序包含分区列每个排序键会映射到一个分区和一个数据文件。该相对权重可以避免为低流量的排序键产生过多小文件。它是一个 double 值定义每个排序键的最小权重。2.0表示每个键具有每个 writer task 目标流量权重2%的基准权重。示例假设 sink Iceberg 表按事件时间每日分区。数据流包含从现在到 180 天前的事件。按事件时间划分时不同日期间的流量权重分布通常呈长尾形态——当天流量最大越老的日期长尾流量越少。假设 writer 并行度为10180 天的总权重为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 optionsIceberg 还提供一组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 configurationDefaultDescriptiontable.exec.iceberg.infer-source-parallelismtrue为 true 时批读取的 source 并行度由 scan split 数量推断得出上限为table.exec.iceberg.infer-source-parallelism.max且受查询 limit若设置约束为 false 时source 并行度取自 Flink 配置。流式读取从不推断并行度。table.exec.iceberg.infer-source-parallelism.max100source 算子推断并行度的最大值。table.exec.iceberg.fetch-batch-record-count2048Iceberg source reader 每次 fetch 批次的目标准入记录数。table.exec.iceberg.worker-pool-sizemax(2, available cpu)用于规划或扫描 manifest 的 worker 池大小。默认为共享 Iceberg worker 池大小受iceberg.worker.num-threads系统属性控制。table.exec.iceberg.use-v2-sinkfalse使用基于 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默认2048FlinkConfigOptions.java控制 source reader 每个 fetch 批次的记录数目标。table.exec.iceberg.worker-pool-size默认取ThreadPools.WORKER_THREAD_POOL_SIZEFlinkConfigOptions.java即共享 Iceberg worker 池大小可通过iceberg.worker.num-threads系统属性调整。table.exec.iceberg.use-v2-sink默认falseFlinkConfigOptions.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 optionDataStream 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 配置中个别作业需要差异化时再用读/写选项按作业覆盖而无需改动表定义或全局配置。赞分享数据湖大数据数据存储【免费下载链接】icebergApache Iceberg项目地址https://gitcode.com/gh_mirrors/icebe/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接口测试上一篇Bytebase React 覆盖层分层策略overlay / agent / critical 三层语义化 z-index 治理与实践下一篇5分钟快速上手C-Qwen3-Embedding-Reranker-0.6B轻量级文本嵌入模型的终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表