Apache Flink

Flink 配置

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

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

Flink 配置

Catalog 配置

通过执行以下语句来创建并命名 catalog(请将 <catalog_name> 替换为你的 catalog 名称,将 <config_key>=<config_value> 替换为 catalog 实现的配置项):

CREATE CATALOG <catalog_name> WITH (
  'type'='iceberg',
  `<config_key>`=`<config_value>`
);

以下属性可以全局设置,不限于某种特定的目录(catalog)实现:

属性是否必填取值说明
type✔️iceberg必须为 iceberg。
catalog-typehive、hadoop、rest、glue、jdbc 或 nessie底层的 Iceberg 目录实现,即 HiveCatalog、HadoopCatalog、RESTCatalog、GlueCatalog、JdbcCatalog、NessieCatalog;若通过 catalog-impl 使用自定义目录实现,则无需设置该属性。
catalog-impl自定义目录实现的全限定类名。如果未设置 catalog-type,则必须设置该属性。
property-version用于描述属性版本的版本号。当属性格式发生变化时,该属性可用于向后兼容。当前的属性版本为 1。
cache-enabledtrue 或 false是否启用目录缓存,默认值为 true。
cache.expiration-interval-ms目录条目在本地缓存的时长,单位为毫秒;-1 等负值将禁用过期,不允许设置为 0。默认值为 -1。

如果使用 Hive 目录,可以设置以下属性:

属性必填值描述
uri✔️Hive metastore 的 thrift URI。
clientsHive metastore 客户端池大小,默认值为 2。
warehouseHive 数据仓库位置。如果既未设置 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 配置值。

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

属性必填取值说明
warehouse✔️存储元数据文件和数据文件的 HDFS 目录。

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

属性必填取值说明
uri✔️REST Catalog 的 URL。
credential在 OAuth2 客户端凭据流程中用于换取令牌的凭据。
token用于与服务端交互的令牌。

运行时配置

读取选项

Flink 读取选项在配置 Flink IcebergSource 时传入:

IcebergSource.forRowData()
    .tableLoader(TableLoader.fromCatalog(...))
    .assignerFactory(new SimpleSplitAssignerFactory())
    .streaming(true)
    .streamingStartingStrategy(StreamingStartingStrategy.INCREMENTAL_FROM_SNAPSHOT_ID)
    .startSnapshotId(3821550127947089987L)
    .monitorInterval(Duration.ofMillis(10L)) // or .set("monitor-interval", "10s") \ set(FlinkReadOptions.MONITOR_INTERVAL, "10s")
    .build()

对于 Flink SQL,读取选项可以通过 SQL hint 传入,如下所示:

SELECT * FROM tableName /*+ OPTIONS('monitor-interval'='10s') */
...

选项可以通过 Flink 配置传入,并将应用于当前会话。请注意,并非所有选项都支持这种方式。

env.getConfig()
    .getConfiguration()
    .set(FlinkReadOptions.SPLIT_FILE_OPEN_COST_OPTION, 1000L);
...

Read option 优先级最高,其次是 Flink configuration,然后是 Table property。

读取选项Flink 配置表属性默认值说明
snapshot-id不适用不适用null用于批处理模式下的时间旅行。从指定的 snapshot-id 读取数据。
case-sensitiveconnector.iceberg.case-sensitive不适用false如果为 true,则以区分大小写的方式匹配列名。
as-of-timestamp不适用不适用null用于批处理模式下的时间旅行。读取给定时间(以毫秒为单位)时点上最近的快照数据。
starting-strategyconnector.iceberg.starting-strategy不适用INCREMENTAL_FROM_LATEST_SNAPSHOT流式执行的启动策略。TABLE_SCAN_THEN_INCREMENTAL:先进行常规表扫描,然后切换到增量模式。增量模式从当前快照的下一个开始。INCREMENTAL_FROM_LATEST_SNAPSHOT:从最新的快照(包含该快照)开始增量模式。如果是空表,则应发现所有后续的追加快照。INCREMENTAL_FROM_LATEST_SNAPSHOT_EXCLUSIVE:从最新的快照(不包含该快照)开始增量模式。如果是空表,则应发现所有后续的追加快照。INCREMENTAL_FROM_EARLIEST_SNAPSHOT:从最早的快照(包含该快照)开始增量模式。如果是空表,则应发现所有后续的追加快照。INCREMENTAL_FROM_SNAPSHOT_ID:从具有特定 id 的快照(包含该快照)开始增量模式。INCREMENTAL_FROM_SNAPSHOT_TIMESTAMP:从具有特定时间戳的快照(包含该快照)开始增量模式。如果该时间戳位于两个快照之间,则应从该时间戳之后的快照开始。仅适用于 FIP27 Source。
start-snapshot-timestamp不适用不适用null从给定时间(以毫秒为单位)时点上最近的快照开始读取数据。
start-snapshot-id不适用不适用null从指定的 snapshot-id 开始读取数据。
end-snapshot-id不适用不适用最新快照 id指定结束快照。
branch不适用不适用main指定批处理模式下读取所用的分支。
tag不适用不适用null指定批处理模式下读取所用的标签。
start-tag不适用不适用null指定增量读取的起始标签。
end-tag不适用不适用null指定增量读取的结束标签。
split-sizeconnector.iceberg.split-sizeread.split.target-size128 MB合并输入分片时的目标大小。
split-lookbackconnector.iceberg.split-file-open-costread.split.planning-lookback10合并输入分片时考虑的分桶数量。
split-file-open-costconnector.iceberg.split-file-open-costread.split.open-file-cost4MB打开文件的预估开销,用作合并分片时的最小权重。
streamingconnector.iceberg.streaming不适用false设置当前任务是以流式还是批处理模式运行。
monitor-intervalconnector.iceberg.monitor-interval不适用60s从新快照中发现分片的监控间隔。仅适用于流式读取。
include-column-statsconnector.iceberg.include-column-stats不适用false基于此创建一个新的扫描,加载每个数据文件的列统计信息。列统计信息包括:值数量、空值数量、下界和上界。
max-planning-snapshot-countconnector.iceberg.max-planning-snapshot-count不适用Integer.MAX_VALUE每次分片枚举所限制的最大快照数量。仅适用于流式读取。
limitconnector.iceberg.limit不适用-1限制输出的行数。
max-allowed-planning-failuresconnector.iceberg.max-allowed-planning-failures不适用3作业失败前允许的扫描计划最大连续失败次数。设置为 -1 表示扫描计划失败时永不使作业失败。
watermark-columnconnector.iceberg.watermark-column不适用null指定用于生成水位线的水位列。如果存在此选项,splitAssignerFactory 将被覆盖为 OrderedSplitAssignerFactory。
watermark-column-time-unitconnector.iceberg.watermark-column-time-unit不适用TimeUnit.MICROSECONDS指定用于生成水位线的水位线时间单位。可选值为 DAYS、HOURS、MINUTES、SECONDS、MILLISECONDS、MICROSECONDS、NANOSECONDS。

写入选项

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 选项默认值说明
write-formatTable write.format.default此次写入操作所使用的文件格式:parquet、avro 或 orc
target-file-size-bytes随表属性而定覆盖此表的 write.target-file-size-bytes
upsert-enabledTable write.upsert.enabled覆盖此表的 write.upsert.enabled
overwrite-enabledfalse覆写表中的数据;当配置为使用 UPSERT 数据流时,不应启用覆写模式。
distribution-modeTable write.distribution-mode覆盖此表的 write.distribution-mode。RANGE 分发模式目前处于实验阶段。
range-distribution-statistics-typeAutoRange 分发的数据统计收集类型:Map、Sketch、Auto。详见此处。
range-distribution-sort-key-base-weight0.0(double)每个排序键相对于每个写入任务目标流量权重的基础权重。详见此处。
compression-codecTable write.(fileformat).compression-codec覆盖此次写入此表所使用的压缩编解码器
compression-levelTable write.(fileformat).compression-level覆盖此次写入 Parquet 和 Avro 表所使用的压缩级别
compression-strategyTable write.orc.compression-strategy覆盖此次写入 ORC 表所使用的压缩策略
write-parallelism上游算子并行度覆盖写入器的并行度
uid-suffix随表属性而定覆盖此表在底层 IcebergSink 中使用的 uid 后缀

Range 分布统计类型

配置值为枚举类型:Map、Sketch、Auto。

  • Map:为每个键收集精确的采样计数。适用于低基数场景(如几百或几千)。
  • Sketch:通过蓄水池采样构建均匀随机采样。非常适合高基数场景(如数百万),因为内存占用保持在较低水平。
  • Auto:以 Map 统计开始。但当检测到基数超过阈值(当前为 10,000)时,统计会自动切换为 Sketch。

Range 分布排序键基础权重

range-distribution-sort-key-base-weight:0.0。

如果排序顺序中包含分区列,则每个排序键都会映射到一个分区和一个数据文件。该相对权重可以避免为流量较低的排序键生成过多的小文件。它是一个双精度浮点值,定义每个排序键的最小权重。0.02 表示每个键的基础权重为目标流量权重的 2%(针对每个写入任务)。

例如,目标 Iceberg 表按事件时间进行每日分区。假设数据流包含从现在到 180 天前的事件。按事件时间来看,不同天数之间的流量权重分布通常呈现长尾模式:当天流量最多,越往后的天数(长尾)流量越少。假设写入并行度为 10,全部 180 天的总权重为 10,000,则每个写入任务的目标流量权重为 1,000。假设最久远的 150 天权重之和为 1,000。通常情况下,range 分区器会把最久远的 150 天全部放入同一个写入任务,该任务将写入 150 个小文件(每天一个)。同时保持 150 个打开的文件可能会消耗大量内存,在检查点时刷写并上传 150 个文件(无论多小)也可能非常缓慢。如果将此配置设置为 0.02,则意味着每个排序键的基础权重为目标权重 1,000 的 2%(针对每个写入任务)。这实际上可以避免在同一个写入任务上放置超过 50 个数据文件(每天一个),无论这些文件有多小。

此配置仅适用于低基数场景下的 StatisticsType.Map。对于 StatisticsType.Sketch 的高基数排序列,它们通常不会被用作分区列,否则在写入时会产生过多的分区和小文件。Sketch range 分区器只是将高基数的键拆分为有序的范围。

评论

登录后参与评论

正在加载评论…