Apache Flink

Flink 写入

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

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

Flink 写入

Iceberg 支持通过 Apache Flink 的 DataStream API 和 Table API 进行批处理与流式写入。

Flink Iceberg Sink 提供精确一次(exactly-once)语义保证。

使用 SQL 写入

Iceberg 同时支持 INSERT INTO 和 INSERT OVERWRITE。

INSERT INTO

要通过 Flink 流式作业向表中追加新数据,请使用 INSERT INTO:

INSERT INTO `hive_catalog`.`default`.`sample` VALUES (1, 'a');
INSERT INTO `hive_catalog`.`default`.`sample` SELECT id, data from other_kafka_table;

INSERT OVERWRITE

要将查询结果替换表中的数据,请在批处理作业中使用 INSERT OVERWRITE(Flink 流式作业不支持 INSERT OVERWRITE)。对于 Iceberg 表而言,覆盖是原子操作。

SELECT 查询产生数据的分区将会被替换,例如:

INSERT OVERWRITE sample VALUES (1, 'a');

Iceberg 还支持通过 select 的值来覆盖指定分区:

INSERT OVERWRITE `hive_catalog`.`default`.`sample` PARTITION(data='a') SELECT 6;

对于分区 Iceberg 表,如果在 PARTITION 子句中为所有分区列都指定了值,则是插入静态分区;如果在 PARTITION 子句中只为部分分区列(即所有分区列的前缀部分)指定了值,则是将查询结果写入动态分区。对于未分区的 Iceberg 表,INSERT OVERWRITE 会完全覆盖其数据。

UPSERT

在向 v2 表格式写入数据时,Iceberg 支持基于主键的 UPSERT。启用 upsert 有两种方式。

  1. 通过表级属性 write.upsert.enabled 启用 UPSERT 模式。以下是创建表时设置该表属性的 SQL 语句示例。除非随后被写入选项覆盖,否则该设置将应用于写入此表的所有路径(批处理或流处理)。

    CREATE TABLE `hive_catalog`.`default`.`sample` (
        `id` INT COMMENT 'unique id',
        `data` STRING NOT NULL,
        PRIMARY KEY(`id`) NOT ENFORCED
    ) with ('format-version'='2', 'write.upsert.enabled'='true');
  2. 使用写入选项中的 upsert-enabled 启用 UPSERT 模式,这比表级配置更灵活。注意,你仍需使用 v2 表格式,并在创建表时指定主键或标识字段。

    INSERT INTO tableName /*+ OPTIONS('upsert-enabled'='true') */
    ...

信息

OVERWRITE 模式与 UPSERT 模式互斥,不能同时启用。在分区表上使用 UPSERT 模式时,等值字段中必须包含对应分区字段的源列。例如,若分区字段为 days(ts),则 ts 必须作为等值字段的一部分。

使用 DataStream 写入

Iceberg 支持从不同的 DataStream 输入写入 Iceberg 表。

追加数据

Flink 原生支持将 DataStream<RowData> 和 DataStream<Row> 写入目标 Iceberg 表。

StreamExecutionEnvironment env = ...;

DataStream<RowData> input = ... ;
Configuration hadoopConf = new Configuration();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path", hadoopConf);

FlinkSink.forRowData(input)
    .tableLoader(tableLoader)
    .append();

env.execute("Test Iceberg DataStream");

覆写数据

在 FlinkSink 构建器中设置 overwrite 标志,即可覆写现有 Iceberg 表中的数据:

StreamExecutionEnvironment env = ...;

DataStream<RowData> input = ... ;
Configuration hadoopConf = new Configuration();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path", hadoopConf);

FlinkSink.forRowData(input)
    .tableLoader(tableLoader)
    .overwrite(true)
    .append();

env.execute("Test Iceberg DataStream");

更新数据

在 FlinkSink 构建器中设置 upsert 标志,即可对现有 Iceberg 表中的数据执行更新操作。该表必须使用 v2 表格式,并且具有主键。

StreamExecutionEnvironment env = ...;

DataStream<RowData> input = ... ;
Configuration hadoopConf = new Configuration();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path", hadoopConf);

FlinkSink.forRowData(input)
    .tableLoader(tableLoader)
    .upsert(true)
    .append();

env.execute("Test Iceberg DataStream");

信息

OVERWRITE 模式与 UPSERT 模式互斥,不能同时启用。在分区表上使用 UPSERT 模式时,必须在相等字段(equality fields)中包含对应分区字段的源列。例如,如果分区字段是 days(ts),那么 ts 必须是相等字段的一部分。

使用 Avro GenericRecord 写入

Flink Iceberg sink 提供了 AvroGenericRecordToRowDataMapper,用于将 Avro GenericRecord 转换为 Flink 的 RowData。你可以使用该 mapper 将 Avro GenericRecord DataStream 写入 Iceberg。

请确保 flink-avro jar 包含在 classpath 中。另外,iceberg-flink-runtime 着色(shaded)打包 jar 无法使用,因为该 runtime jar 对 avro 包进行了着色处理。请改用非着色的 iceberg-flink jar。

DataStream<org.apache.avro.generic.GenericRecord> dataStream = ...;

Schema icebergSchema = table.schema();


// The Avro schema converted from Iceberg schema can't be used
// due to precision difference between how Iceberg schema (micro)
// and Flink AvroToRowDataConverters (milli) deal with time type.
// Instead, use the Avro schema defined directly.
// See AvroGenericRecordToRowDataMapper Javadoc for more details.
org.apache.avro.Schema avroSchema = AvroSchemaUtil.convert(icebergSchema, table.name());

GenericRecordAvroTypeInfo avroTypeInfo = new GenericRecordAvroTypeInfo(avroSchema);
RowType rowType = FlinkSchemaUtil.convert(icebergSchema);

FlinkSink.builderFor(
    dataStream,
    AvroGenericRecordToRowDataMapper.forAvroSchema(avroSchema),
    FlinkCompatibilityUtil.toTypeInfo(rowType))
  .table(table)
  .tableLoader(tableLoader)
  .append();

分支写入

Iceberg 表也支持通过 FlinkSink 中的 toBranch API 写入分支。有关分支的更多信息,请参阅分支。

FlinkSink.forRowData(input)
    .tableLoader(tableLoader)
    .toBranch("audit-branch")
    .append();

指标

Flink Iceberg sink 提供以下 Flink 指标。

并行写入器的指标位于 IcebergStreamWriter 子分组下,应包含以下键值标签。

  • table:完整表名(如 iceberg.my_db.my_table)
  • subtask_index:写入器子任务索引,从 0 开始
指标名称指标类型描述
lastFlushDurationMsGauge写入子任务在检查点期间刷新并上传文件所花费的时长(毫秒)。
flushedDataFilesCounter已刷新并上传的数据文件数量。
flushedDeleteFilesCounter已刷新并上传的删除文件数量。
flushedReferencedDataFilesCounter已刷新的删除文件所引用的数据文件数量。
dataFilesSizeHistogramHistogram数据文件大小(字节)的直方图分布。
deleteFilesSizeHistogramHistogram删除文件大小(字节)的直方图分布。

上述 Histogram 指标需要 classpath 中存在 org.apache.flink:flink-metrics-dropwizard,而 Flink 默认并不附带该依赖。请将此构件添加到你的 classpath 中,才能看到直方图指标。如果未添加,直方图指标将不会出现,其他所有类型的指标仍会正常发布。

提交器(Committer)的指标位于 IcebergFilesCommitter 子分组下,应包含以下键值标签。

  • table:完整表名(如 iceberg.my_db.my_table)
指标名称指标类型描述
lastCheckpointDurationMsGauge提交器算子对其状态执行检查点所耗费的时长(毫秒)。
lastCommitDurationMsGaugeIceberg 表提交所耗费的时长(毫秒)。
committedDataFilesCountCounter已提交的数据文件数量。
committedDataFilesRecordCountCounter已提交的数据文件中包含的记录数。
committedDataFilesByteCountCounter已提交的数据文件中包含的字节数。
committedDeleteFilesCountCounter已提交的删除文件数量。
committedDeleteFilesRecordCountCounter已提交的删除文件中包含的记录数。
committedDeleteFilesByteCountCounter已提交的删除文件中包含的字节数。
elapsedSecondsSinceLastSuccessfulCommitGauge距离上一次成功 Iceberg 提交所经过的时间(秒)。

elapsedSecondsSinceLastSuccessfulCommit 是用于检测 Iceberg 提交失败或缺失的理想告警指标。

  • Iceberg 提交发生在 Flink 检查点成功之后的 notifyCheckpointComplete 回调中。可能出现这样的情况:Flink 检查点成功,而 Iceberg 提交由于某种原因失败了。
  • 也可能出现 notifyCheckpointComplete 未被触发的情况(无论出于什么缺陷)。其结果是,根本不会有 Iceberg 提交被尝试执行。

如果检查点间隔(即预期的 Iceberg 提交间隔)为 5 分钟,请设置类似 elapsedSecondsSinceLastSuccessfulCommit > 60 minutes 规则的告警,以检测过去一小时内失败或缺失的 Iceberg 提交。

选项

写入选项

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') */
...

在此查看所有可用选项:write-options

分发模式

Flink 流式写入器同时支持 HASH 和 RANGE 分发模式。可以通过 FlinkSink#Builder#distributionMode(DistributionMode) 启用,也可以通过 write-options 启用。

哈希分发

HASH 分发根据分区键(分区表)或等值字段(非分区表)对数据进行重分布,它直接利用 Flink 的 DataStream#keyBy 来分发数据。

HASH 分发有一些局限性:

  • 它无法很好地处理数据倾斜。例如,某些分区的数据量远大于其他分区。
  • 如果分区键或等值字段的基数较低,可能导致流量分布不均衡,如 [PR 4228](https://github.com/apache/iceberg/pull/4228) 中所示。
  • 写入器的并行度受限于哈希键的基数。如果基数为 10,则最多只有 10 个写入任务能够接收到流量。即使流量规模需要更高的写入并行度,提高并行度也不会带来帮助。

范围分发(实验性)

RANGE 分发通过自定义范围分区器,按照分区键或排序顺序对数据进行重分布。范围分发会收集流量统计数据,以引导范围分区器将流量均匀地分配给写入任务。

范围分发仅通过范围分区器对数据进行重分布。数据文件中的行不会被排序,因为 Flink 流式写入器目前尚不支持这一点。

使用场景

范围分发适用于已分区或定义了 SortOrder 的 Iceberg 表。对于已分区但未定义 SortOrder 的表,会使用分区列作为排序顺序。如果表显式定义了 SortOrder,范围分区器将使用该排序顺序。

范围分发能够处理数据倾斜。例如:

  • 表按事件时间分区。通常最近几小时的数据较多,而长尾时段的数据越来越少。
  • 表按国家/地区代码分区,其中某些国家(如美国)的流量远大于其他国家,而较小国家的数据量少得多。
  • 表按事件类型分区,其中某些类型的数据量远大于其他类型。

范围分发还可以在非分区列上对数据进行聚类。例如,表按摄取时间进行小时级分区,而查询通常包含针对非分区列(如 device_id 或 country_code)的谓词。当表的 SortOrder 定义中包含非分区列时,范围分区可以通过在该非分区列上进行聚类来提升查询性能。

流量统计数据

统计信息由每个 shuffle 算子的子任务收集,并由协调器在每个检查点周期内进行聚合。聚合后的统计信息会广播给所有子任务,并在下一个检查点时应用到范围分区器。因此,检测到流量分布变化并将新统计信息应用到范围分区器,可能最多需要两个检查点周期。

范围分布既可用于低基数场景(如 country_code),也可用于高基数场景(如 device_id)。

  • 对于低基数场景(如几百或几千个取值),使用 HashMap 跟踪每个键的流量分布。如果出现新的排序键值,在学到该新键的流量分布之前,范围分区器会先以轮询方式将其分配给写入任务。
  • 对于高基数场景(如数百万或数十亿个取值),使用均匀随机采样(蓄水池采样)计算能够均分排序键空间的范围边界。这种方式可以保持较低的内存占用和网络交换量。当键的分布相对均匀时,蓄水池采样效果良好。如果某个热点键占据了明显过大的流量份额,通过均匀采样进行的范围划分可能效果不佳。

使用方式

以下是在 Java 中启用范围分布的方法。有两个可选的高级配置,默认值在大多数场景下都能很好地工作。详情参见 write-options。

FlinkSink.forRowData(input)
    ...
    .distributionMode(DistributionMode.RANGE)
    .rangeDistributionStatisticsType(StatisticsType.Auto)
    .rangeDistributionSortKeyBaseWeight(0.0d)
    .append();

开销

数据洗牌(哈希或范围)在序列化/反序列化和网络 I/O 方面存在计算开销。预计 CPU 使用率会有所上升。

范围分发(Range distribution)还会收集并聚合数据分布统计信息,这同样会带来一定的 CPU 开销。若使用默认的 Auto 统计类型,内存开销通常很小。如果键的基数很高,请不要使用 Map 统计类型,否则可能会导致显著的内存占用,以及统计聚合时大量的网络数据交换。

注意事项

Flink 流式写入作业依赖快照摘要来记录最近一次提交的检查点 ID,并将未提交的数据存储为临时文件。因此,过期快照(expiring snapshots)和删除孤立文件(deleting orphan files)可能会破坏 Flink 作业的状态。为避免这种情况,请确保保留 Flink 作业创建的最后一个快照(可通过摘要中的 flink.job-id 属性来识别),并且只删除足够旧的孤立文件。

基于 Sink V2 的实现

在当前默认的 FlinkSink 实现创建之时,Flink Sink 的接口存在一些无法满足 Iceberg 表需求的局限性。由于这些局限性,FlinkSink 基于一个自定义的 StreamOperator 链构建,并以 DiscardingSink 结尾。

在 Flink 1.15 版本中引入了 SinkV2 接口。新的 IcebergSink 实现(位于 iceberg-flink 模块中)便使用了该接口。该新实现是后续开展诸如表维护等功能工作的基础。基于 SinkV2 的实现目前是一项实验性功能,使用时请谨慎。

使用 SQL 写入

要在 SQL 中启用基于 SinkV2 的实现,请设置以下配置项:

SET table.exec.iceberg.use-v2-sink = true;

使用 DataStream 写入

要使用基于 SinkV2 的实现,请将代码片段中的 FlinkSink 替换为 IcebergSink。

警告

这两种实现之间存在一些细微差别:

  • IcebergSink 暂不支持 RANGE 分发模式
  • 使用 IcebergSink 时请使用 uidSuffix 而非 uidPrefix

Flink 动态 Iceberg Sink

Flink 动态 Iceberg Sink(Dynamic Sink)支持以下功能:

  1. 写入任意数量的表
    单个 sink 可以动态将记录路由到多个 Iceberg 表。
  2. 动态创建和更新表
    根据用户定义的路由逻辑创建和更新表。
  3. 动态 Schema 和分区演进
    表的 Schema 和分区规范可在流式执行过程中更新。

所有配置均通过 DynamicRecord 类进行控制,因此在需求变更时无需重启 Flink 作业。

    DynamicIcebergSink.forInput(dataStream)
        .generator((inputRecord, out) -> out.collect(
                new DynamicRecord(
                        TableIdentifier.of("db", "table"),
                        "branch",
                        SCHEMA,
                        (RowData) inputRecord,
                        PartitionSpec.unpartitioned(),
                        DistributionMode.HASH,
                        2)))
        .catalogLoader(CatalogLoader.hive("hive", new Configuration(), Map.of()))
        .writeParallelism(10)
        .immediateTableUpdate(true)
        .append();

配置示例

DynamicIcebergSink.Builder<RowData> builder = DynamicIcebergSink.forInput(inputStream);

// Set common properties
builder
    .set("write.parquet.compression-codec", "gzip");

// Set Dynamic Sink specific options
builder
    .writeParallelism(4)
    .uidPrefix("dynamic-sink")
    .cacheMaxSize(500)
    .cacheRefreshMs(5000);

// Add generator and append sink
builder.generator(new CustomRecordGenerator());
builder.append();

动态路由配置

通过实现 DynamicRecordGenerator 接口,可以自定义动态表路由:

public class CustomRecordGenerator implements DynamicRecordGenerator<RowData> {
    @Override
    public DynamicRecord generate(RowData row) {
        DynamicRecord record = new DynamicRecord();
        // Set table name based on business logic
        TableIdentifier tableIdentifier = TableIdentifier.of(database, tableName);
        record.setTableIdentifier(tableIdentifier);
        record.setData(row);
        // Set the maximum number of parallel writers for a given table/branch/schema/spec
        record.writeParallelism(2);
        return record;
    }
}

// Set custom record generator when building the sink
DynamicIcebergSink.Builder<RowData> builder = DynamicIcebergSink.forInput(inputStream);
builder.generator(new CustomRecordGenerator());
// ... other config ...
builder.append();

用户需要提供一个转换器,将输入记录转换为 DynamicRecord。对于每条记录,我们需要以下 DynamicRecord 信息:

属性描述
TableIdentifier该记录将被写入的目标表。
Branch写入该记录的目标分支(可选)。
Schema该记录的 schema。
Spec该记录预期的分区规范。
RowData要写入的实际行数据。
DistributionMode写入该记录的分发模式(NONE、HASH 或 null)。当为 null 时,该记录完全不会被 shuffle。
Parallelism给定表/分支/schema/分区规范(WriteTarget)的最大并行写入者数量。
UpsertMode覆盖该表的 write.upsert.enabled 配置(可选)。
EqualityFields该表的相等性字段(可选)。

Schema 演进

动态 sink 会尝试将 DynamicRecord 中提供的 schema 与现有的表 schema 进行匹配。

  • 如果与某个现有表 schema 直接匹配,则使用该表 schema 写入表。
  • 如果没有直接匹配,DynamicSink 会尝试对提供的 schema 进行适配,使其与某个表 schema 匹配。例如,如果表 schema 中多出一个可选列,则会向通过 DynamicRecord 提供的 RowData 中添加一个 null 值。
  • 否则,我们会在下述约束范围内对表 schema 进行演进,以使其与输入 schema 匹配。

动态 sink 为表元数据和传入的 schema 维护了一个 LRU 缓存,并基于大小和时间约束进行淘汰。当 DynamicRecord 包含的 schema 与当前表 schema 不兼容时,会触发一次 schema 更新。根据 immediateTableUpdate 配置,该更新可以立即执行,也可以通过集中式执行器执行。虽然集中式更新减轻了 Catalog 的负载,但可能会给 sink 带来反压。

支持的 schema 更新

  • 添加新列
  • 扩展已有列的类型(例如 Integer → Long、Float → Double)
  • 将必填列改为可选列
  • 删除列(默认禁用)

删除列默认被禁用,以避免因迟到或乱序数据引发问题,因为被删除的字段一旦移除,就无法在不丢失数据的情况下轻松恢复。

你可以选择启用删除列功能(见下方的配置选项)。列被删除后,从技术上讲仍然可以向该列写入数据,因为 Iceberg 会保留表的所有历史模式。然而,常规查询将无法引用该列。如果该字段作为新模式的一部分重新出现,系统会添加一个全新的列,这个新列除了列名之外与旧列毫无关联,也就是说,对新列的查询永远不会返回旧列的数据。

不支持的模式更新
  • 重命名列

不支持重命名,因为模式比较是基于名称的,而重命名需要额外的元数据或提示才能解析。

缓存

涉及两种不同的缓存:表元数据缓存和输入模式缓存。

  • 表元数据缓存保存模式定义、分区规范等元数据,以减少重复的 Catalog 查询。其大小由 cacheMaxSize 设置控制。
  • 输入模式缓存按表存储传入的模式及其兼容性解析结果。其大小由 inputSchemasPerTableCacheMaxSize 控制。

为提高缓存命中率和性能,如果记录的模式未发生变化,请复用同一个 DynamicRecord.schema 实例。

动态 Sink 配置

Dynamic Iceberg Flink Sink 使用 Builder 模式进行配置。以下是关键的配置方法:

方法说明
overwrite(boolean enabled)启用覆盖模式
writeParallelism(int parallelism)设置写入并行度
uidPrefix(String prefix)设置算子 UID 前缀
snapshotProperties(Map<String, String> properties)设置快照元数据属性
toBranch(String branch)写入到指定分支
cacheMaxSize(int maxSize)设置表元数据的缓存大小
cacheRefreshMs(long refreshMs)设置缓存刷新间隔
inputSchemasPerTableCacheMaxSize(int size)设置每张表缓存的输入 schema 的最大数量
immediateTableUpdate(boolean enabled)控制表元数据(schema/分区规范)是否立即更新(默认:false)
set(String property, String value)设置任意 Iceberg 写入属性(例如 "write.format"、"write.upsert.enabled")。查看此处的全部选项:write-options
setAll(Map<String, String> properties)一次设置多个属性
tableCreator(TableCreator creator)当 DynamicIcebergSink 创建新的 Iceberg 表时,允许覆盖表的创建方式——根据表名设置自定义表属性和存储位置。
dropUnusedColumns(boolean enabled)启用后,会从当前表 schema 中删除输入 schema 未包含的所有列(关于删除列的注意事项,参见上文)。

分发模式

设置在每个 DynamicRecord 上的 DistributionMode 控制记录如何从处理器路由到写入器:

模式行为
NONE记录以轮询方式分发到各个写入器子任务(若设置了相等字段,则按相等字段分发)。
HASH记录按分区键(分区表)或相等字段(非分区表)分发,确保同一分区的记录由同一个写入器子任务处理。
null转发模式:完全跳过分发,通过转发边直接发送记录(见下文)。

转发模式

使用不带 distributionMode 参数的 DynamicRecord 构造函数重载,可完全跳过分发。该模式专为高吞吐管道设计:在这类管道中,每个分区本身就已拥有大量数据,序列化和网络洗牌(shuffle)的开销过高。记录通过转发边从处理器直接发送到写入器,从而启用 Flink 算子链(operator chaining)。表元数据的更新始终在处理器内部立即执行(无论 immediateTableUpdate 设置如何),这是因为刻意省略了专用的表更新算子,以避免引入额外的数据洗牌。

转发记录与常规记录可以在同一条管道中混合使用。处理器会将记录路由到两个独立的 sink 输出:

  • Shuffle sink:接收需要洗牌的记录。这些记录在到达写入器之前,会先经过正常的分发拓扑(哈希/轮询)。
  • Forward sink:接收没有 distributionMode 的记录。这些记录完全跳过分发,通过转发边从处理器直接流动,从而支持 Flink 算子链。适用于避免洗牌开销至关重要的高吞吐表。该 sink 的 writeParallelism 配置不适用于此路径。

警告

  1. 在 forward 路径中,架构变更总是立即生效,因为记录必须通过 forward 边直接传递。对于预期的大批量使用场景,这可能导致向 Iceberg catalog 提交大量冲突,并暂时延迟数据处理。建议在使用新架构发布记录之前先在外部更新架构,或者在上游引入新架构时预留吞吐量的临时中断时间。
  2. 由于 forward 路径完全跳过分发,用户需要在上游自行正确分发数据,再让记录到达动态 Iceberg sink。否则,写入可能会不均衡。

注意事项

  • 范围分发模式:目前动态 sink 不支持 RANGE 分发模式,若设置则会回退为 HASH。
  • 属性优先级说明:当表属性与 sink 属性发生冲突时,sink 属性会覆盖表属性的配置。
  • 表格式版本升级:动态 sink 不支持对包含动态记录的表进行升级。在 V2 到 V3 升级过程中,作业不应处于运行状态。

评论

登录后参与评论

正在加载评论…