Flink 配置
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-type | hive、hadoop、rest、glue、jdbc 或 nessie | 底层的 Iceberg 目录实现,即 HiveCatalog、HadoopCatalog、RESTCatalog、GlueCatalog、JdbcCatalog、NessieCatalog;若通过 catalog-impl 使用自定义目录实现,则无需设置该属性。 | |
| catalog-impl | 自定义目录实现的全限定类名。如果未设置 catalog-type,则必须设置该属性。 | ||
| property-version | 用于描述属性版本的版本号。当属性格式发生变化时,该属性可用于向后兼容。当前的属性版本为 1。 | ||
| cache-enabled | true 或 false | 是否启用目录缓存,默认值为 true。 | |
| cache.expiration-interval-ms | 目录条目在本地缓存的时长,单位为毫秒;-1 等负值将禁用过期,不允许设置为 0。默认值为 -1。 |
如果使用 Hive 目录,可以设置以下属性:
| 属性 | 必填 | 值 | 描述 |
|---|---|---|---|
| uri | ✔️ | Hive metastore 的 thrift URI。 | |
| clients | Hive metastore 客户端池大小,默认值为 2。 | ||
| warehouse | Hive 数据仓库位置。如果既未设置 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-sensitive | connector.iceberg.case-sensitive | 不适用 | false | 如果为 true,则以区分大小写的方式匹配列名。 |
| as-of-timestamp | 不适用 | 不适用 | null | 用于批处理模式下的时间旅行。读取给定时间(以毫秒为单位)时点上最近的快照数据。 |
| starting-strategy | connector.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-size | connector.iceberg.split-size | read.split.target-size | 128 MB | 合并输入分片时的目标大小。 |
| split-lookback | connector.iceberg.split-file-open-cost | read.split.planning-lookback | 10 | 合并输入分片时考虑的分桶数量。 |
| split-file-open-cost | connector.iceberg.split-file-open-cost | read.split.open-file-cost | 4MB | 打开文件的预估开销,用作合并分片时的最小权重。 |
| streaming | connector.iceberg.streaming | 不适用 | false | 设置当前任务是以流式还是批处理模式运行。 |
| monitor-interval | connector.iceberg.monitor-interval | 不适用 | 60s | 从新快照中发现分片的监控间隔。仅适用于流式读取。 |
| include-column-stats | connector.iceberg.include-column-stats | 不适用 | false | 基于此创建一个新的扫描,加载每个数据文件的列统计信息。列统计信息包括:值数量、空值数量、下界和上界。 |
| max-planning-snapshot-count | connector.iceberg.max-planning-snapshot-count | 不适用 | Integer.MAX_VALUE | 每次分片枚举所限制的最大快照数量。仅适用于流式读取。 |
| limit | connector.iceberg.limit | 不适用 | -1 | 限制输出的行数。 |
| max-allowed-planning-failures | connector.iceberg.max-allowed-planning-failures | 不适用 | 3 | 作业失败前允许的扫描计划最大连续失败次数。设置为 -1 表示扫描计划失败时永不使作业失败。 |
| watermark-column | connector.iceberg.watermark-column | 不适用 | null | 指定用于生成水位线的水位列。如果存在此选项,splitAssignerFactory 将被覆盖为 OrderedSplitAssignerFactory。 |
| watermark-column-time-unit | connector.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-format | Table write.format.default | 此次写入操作所使用的文件格式:parquet、avro 或 orc |
| target-file-size-bytes | 随表属性而定 | 覆盖此表的 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-sort-key-base-weight | 0.0(double) | 每个排序键相对于每个写入任务目标流量权重的基础权重。详见此处。 |
| 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 | 上游算子并行度 | 覆盖写入器的并行度 |
| 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 分区器只是将高基数的键拆分为有序的范围。
评论
登录后参与评论
KnowForge