Apache Flink

Flink 表维护

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

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

Flink TableMaintenance

Flink 表维护 BatchMode

重写文件操作

Iceberg 提供了 API,可通过提交 Flink 批处理作业将小文件重写为大文件。此 Flink 操作的行为与 Spark 的 rewriteDataFiles 相同。

import org.apache.iceberg.flink.actions.Actions;

TableLoader tableLoader = TableLoader.fromCatalog(
    CatalogLoader.hive("my_catalog", configuration, properties),
    TableIdentifier.of("database", "table")
);

Table table = tableLoader.loadTable();
RewriteDataFilesActionResult result = Actions.forTable(table)
        .rewriteDataFiles()
        .execute();

有关 rewrite files 操作的更多详细信息,请参阅 RewriteDataFilesAction

Flink Table Maintenance StreamingMode

概述

在 Flink 流处理环境中的 Apache Iceberg 部署里,实现自动化的表维护操作——包括 snapshot expiration(快照过期)、small file compaction(小文件合并)以及 orphan file cleanup(孤立文件清理)——对于获得最佳查询性能和存储效率至关重要。

以往,这些维护操作只能通过 Iceberg Spark Actions 来执行,因此必须部署和管理专用的 Spark 集群。仅仅为了优化表而依赖 Spark 基础设施,会带来显著的架构复杂性和运维开销。

Apache Iceberg 中的 TableMaintenance API 使 Flink 作业能够原生执行维护任务,既可嵌入现有的流处理管道中,也可作为独立的 Flink 作业部署。这消除了对外部系统的依赖,从而简化架构、降低运维成本并增强自动化能力。

支持的特性(Flink)

ExpireSnapshots

移除过期的快照及其对应文件。提交时内部会使用 cleanExpiredFiles(true),因此过期的元数据和文件会被自动清理。

.add(ExpireSnapshots.builder()
    .maxSnapshotAge(Duration.ofDays(7))
    .retainLast(10)
    .deleteBatchSize(1000))

RewriteDataFiles

压缩小文件以优化文件大小。支持部分进度提交,并可限制每次运行重写的最大字节数。

.add(RewriteDataFiles.builder()
    .targetFileSizeBytes(256 * 1024 * 1024)
    .minFileSizeBytes(32 * 1024 * 1024)
    .partialProgressEnabled(true)
    .partialProgressMaxCommits(5))

DeleteOrphanFiles

用于删除未被 Iceberg 表的任何元数据文件所引用的文件,这类文件可被视为"孤立文件"(orphaned)。该操作会检查表的位置目录,找出此类文件。

.add(DeleteOrphanFiles.builder()
    .minAge(Duration.ofDays(3))
    .deleteBatchSize(1000))

锁管理

TriggerLockFactory 对于协调维护任务至关重要。它可防止对同一张表执行并发维护操作,从而避免可能产生的冲突或数据损坏。即使对于单个作业,这种锁定机制也是必要的,因为同一任务的多个实例也可能发生冲突。

为什么需要锁

  • 并发访问:多个 Flink 作业可能同时尝试执行维护操作
  • 数据一致性:确保同一张表在同一时间只运行一个维护操作
  • 资源管理:防止资源冲突和调度问题
  • 避免重复工作:即使只调度了一个 compaction 作业,多个实例也可能尝试执行相同的操作,从而导致冗余工作和资源浪费。

支持的锁类型

JDBC 锁工厂

使用数据库表来管理分布式锁:

Map<String, String> jdbcProps = new HashMap<>();
jdbcProps.put("jdbc.user", "flink");
jdbcProps.put("jdbc.password", "flinkpw");
jdbcProps.put("flink-maintenance.lock.jdbc.init-lock-tables", "true"); // Auto-create lock table if it doesn't exist

TriggerLockFactory lockFactory = new JdbcLockFactory(
    "jdbc:postgresql://localhost:5432/iceberg", // JDBC URL
    "catalog.db.table",                         // Lock ID (unique identifier)
    jdbcProps                                   // JDBC connection properties
);
ZooKeeper 锁工厂

使用 Apache ZooKeeper 实现分布式锁:

TriggerLockFactory lockFactory = new ZkLockFactory(
    "localhost:2181",       // ZooKeeper connection string
    "catalog.db.table",     // Lock ID (unique identifier)
    60000,                  // sessionTimeoutMs
    15000,                  // connectionTimeoutMs
    3000,                   // baseSleepTimeMs
    3                       // maxRetries
);

Flink 维护的锁

在 Flink 内部维护该锁。这不需要配置外部系统,唯一的前提是同一张表不能存在并行执行的表维护作业。

快速开始

下面的示例演示了如何在 Flink 环境中实现 Iceberg 表的自动维护。

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

TableLoader tableLoader = TableLoader.fromCatalog(
    CatalogLoader.hive("my_catalog", configuration, properties),
    TableIdentifier.of("database", "table")
);

Map<String, String> jdbcProps = new HashMap<>();
jdbcProps.put("jdbc.user", "flink");
jdbcProps.put("jdbc.password", "flinkpw");

// JdbcLockFactory Example
TriggerLockFactory lockFactory = new JdbcLockFactory(
    "jdbc:postgresql://localhost:5432/iceberg", // JDBC URL
    "catalog.db.table",                         // Lock ID (unique identifier)
    jdbcProps                                   // JDBC connection properties
);

// Option 1: With external lock factory (plan to deprecate this Option since 1.12)
TableMaintenance.forTable(env, tableLoader, lockFactory)
// Option 2: With Flink-managed lock (no external lock required)
TableMaintenance.forTable(env, tableLoader)
    .uidSuffix("my-maintenance-job")
    .rateLimit(Duration.ofMinutes(10))
    .lockCheckDelay(Duration.ofSeconds(10))
    .add(ExpireSnapshots.builder()
        .scheduleOnCommitCount(10)
        .maxSnapshotAge(Duration.ofMinutes(10))
        .retainLast(5)
        .deleteBatchSize(5)
        .parallelism(8))
    .add(RewriteDataFiles.builder()
        .scheduleOnDataFileCount(10)
        .targetFileSizeBytes(128 * 1024 * 1024)
        .partialProgressEnabled(true)
        .partialProgressMaxCommits(10))
    .append();

env.execute("Table Maintenance Job");

配置选项

TableMaintenance 构建器

方法说明默认值
uidSuffix(String)作业的唯一标识符后缀随机 UUID
rateLimit(Duration)两次任务执行之间的最小间隔60 秒
lockCheckDelay(Duration)检查锁可用性的延迟30 秒
parallelism(int)维护任务的默认并行度系统默认值
maxReadBack(int)初始化期间要检查的最大快照数100

维护任务通用选项

方法说明默认值类型
scheduleOnCommitCount(int)提交 N 次后触发无自动调度int
scheduleOnDataFileCount(int)产生 N 个数据文件后触发无自动调度int
scheduleOnDataFileSize(long)数据文件总大小(字节)达到后触发无自动调度long
scheduleOnPosDeleteFileCount(int)产生 N 个 position delete 文件后触发无自动调度int
scheduleOnPosDeleteRecordCount(long)产生 N 条 position delete 记录后触发无自动调度long
scheduleOnEqDeleteFileCount(int)产生 N 个 equality delete 文件后触发无自动调度int
scheduleOnEqDeleteRecordCount(long)产生 N 条 equality delete 记录后触发无自动调度long
scheduleOnInterval(Duration)时间间隔达到后触发无自动调度Duration

ExpireSnapshots 配置

方法说明默认值类型
maxSnapshotAge(Duration)要保留的快照的最大存活时长5 天Duration
retainLast(int)要保留的最少快照数量1int
deleteBatchSize(int)每批次删除的文件数量1000int
planningWorkerPoolSize(int)用于规划快照过期的工作线程数量共享工作线程池int
cleanExpiredMetadata(boolean)过期快照时删除过期的元数据文件trueboolean

RewriteDataFiles 配置

方法说明默认值类型
targetFileSizeBytes(long)重写文件的目标大小表属性或 512MBlong
minFileSizeBytes(long)可参与压缩的文件的最小大小目标文件大小的 75%long
maxFileSizeBytes(long)可参与压缩的文件的最大大小目标文件大小的 180%long
minInputFiles(int)触发重写的最少文件数5int
deleteFileThreshold(int)每个数据文件强制重写所需的最小删除文件数Integer.MAX_VALUEint
rewriteAll(boolean)无论阈值如何,重写所有数据文件falseboolean
maxFileGroupSizeBytes(long)文件组的最大总大小107374182400(100GB)long
maxFilesToRewrite(int)若未指定此选项,将重写所有符合条件的文件nullint
partialProgressEnabled(boolean)启用部分进度提交falseboolean
partialProgressMaxCommits(int)partialProgressEnabled 为 true 时,部分进度允许的最大提交次数10int
maxRewriteBytes(long)每次执行允许重写的最大字节数Long.MAX_VALUElong
filter(Expression)用于选择待重写文件的过滤表达式Expressions.alwaysTrue()Expression
maxFileGroupInputFiles(long)文件组内允许的最大输入文件数Long.MAX_VALUElong

DeleteOrphanFiles 配置

方法说明默认值类型
location(string)开始递归列出待移除候选文件的位置表的位置String

usePrefixListing(boolean) 设为 true 时,通过 SupportsPrefixOperations 接口使用基于前缀的文件列表。启用此标志时,Table 的 FileIO 实现必须支持 SupportsPrefixOperations。(注意:设为 false 将使用递归方法获取文件信息。如果底层存储是对象存储,则会反复调用 API 来获取路径。) 默认值:true 类型:boolean

prefixMismatchMode(PrefixMismatchMode) 当位置前缀(scheme/authority)不匹配时的操作行为:

  • ERROR - 抛出异常。
  • IGNORE - 不做任何操作。
  • DELETE - 删除文件。

默认值:ERROR 类型:PrefixMismatchMode

equalSchemes(Map<String, String>) 需要视为相等的文件系统 scheme 映射。键为以逗号分隔的 scheme 列表,值为单个 scheme。示例:"s3n"=>"s3","s3a"=>"s3" 默认值:Map 类型:Map

equalAuthorities(Map<String, String>) 需要视为相等的文件系统 authority 映射。键为以逗号分隔的 authority 列表,值为单个 authority。默认值:空 Map 类型:Map

minAge(Duration) 删除在该时间戳之前创建的孤儿文件 默认值:3 天前 类型:Duration

planningWorkerPoolSize(int) 用于规划快照过期的工作线程数 默认值:共享工作线程池 类型:int

完整示例

public class TableMaintenanceJob {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(60000); // Enable checkpointing

        // Configure table loader
        TableLoader tableLoader = TableLoader.fromCatalog(
            CatalogLoader.hive("my_catalog", configuration),
            TableIdentifier.of("database", "table")
        );

        // Set up JDBC lock factory
        Map<String, String> jdbcProps = new HashMap<>();
        jdbcProps.put("jdbc.user", "flink");
        jdbcProps.put("jdbc.password", "flinkpw");
        jdbcProps.put("flink-maintenance.lock.jdbc.init-lock-tables", "true");

        TriggerLockFactory lockFactory = new JdbcLockFactory(
            "jdbc:postgresql://localhost:5432/iceberg",
            "catalog.db.table",
            jdbcProps
        );

        // Set up maintenance with comprehensive configuration
        TableMaintenance.forTable(env, tableLoader, lockFactory)
            .uidSuffix("production-maintenance")
            .rateLimit(Duration.ofMinutes(15))
            .lockCheckDelay(Duration.ofSeconds(30))
            .parallelism(4)

            // Daily snapshot cleanup
            .add(ExpireSnapshots.builder()
                .maxSnapshotAge(Duration.ofDays(7))
                .retainLast(10))

            // Continuous file optimization
            .add(RewriteDataFiles.builder()
                .targetFileSizeBytes(256 * 1024 * 1024)
                .minFileSizeBytes(32 * 1024 * 1024)
                .scheduleOnDataFileCount(20)
                .partialProgressEnabled(true)
                .partialProgressMaxCommits(5)
                .maxRewriteBytes(2L * 1024 * 1024 * 1024)
                .parallelism(6))

            // Delete orphans files created more than five days ago
            .add(DeleteOrphanFiles.builder()
                        .minAge(Duration.ofDays(5)))

            .append();

        env.execute("Iceberg Table Maintenance");
    }
}

IcebergSink 与提交后任务集成

Apache Iceberg Sink V2 为 Flink 提供了在数据提交到表之后自动执行维护任务的能力,通过 addPostCommitTopology(...) 方法实现。

DataStream API

Builder
IcebergSink.forRowData(dataStream)
    .table(table)
    .tableLoader(tableLoader)
    .rewriteDataFiles(Map.of(
        RewriteDataFilesConfig.MAX_BYTES, "1073741824"))
    .expireSnapshots(Map.of(
        ExpireSnapshotsConfig.RETAIN_LAST, "5",
        ExpireSnapshotsConfig.MAX_SNAPSHOT_AGE_SECONDS, "604800"))
    .deleteOrphanFiles(Map.of(
        DeleteOrphanFilesConfig.MIN_AGE_SECONDS, "259200"))
    .append();
配置

所有维护任务均通过字符串属性进行配置:

Map<String, String> flinkConf = new HashMap<>();

// Enable maintenance tasks
flinkConf.put("flink-maintenance.rewrite.enabled", "true");
flinkConf.put("flink-maintenance.expire-snapshots.enabled", "true");
flinkConf.put("flink-maintenance.delete-orphan-files.enabled", "true");

// Configure rewrite data files
flinkConf.put("flink-maintenance.rewrite.max-bytes", "1073741824");

// Configure expire snapshots
flinkConf.put("flink-maintenance.expire-snapshots.retain-last", "5");
flinkConf.put("flink-maintenance.expire-snapshots.max-snapshot-age-seconds", "604800");

// Configure delete orphan files
flinkConf.put("flink-maintenance.delete-orphan-files.min-age-seconds", "259200");

// Configure JDBC lock settings (deprecated, lock configuration is no longer required for a single Flink job)
flinkConf.put("flink-maintenance.lock.type", "jdbc");
flinkConf.put("flink-maintenance.lock.jdbc.uri", "jdbc:postgresql://localhost:5432/iceberg");
flinkConf.put("flink-maintenance.lock.lock-id", "catalog.db.table");

IcebergSink.forRowData(dataStream)
    .table(table)
    .tableLoader(tableLoader)
    .setAll(flinkConf)
    .append();

SQL 示例

你可以在执行写入操作之前,通过 SQL 启用维护并配置锁:

-- Enable Iceberg V2 Sink and maintenance tasks
SET 'table.exec.iceberg.use.v2.sink' = 'true';
SET 'flink-maintenance.rewrite.enabled' = 'true';
SET 'flink-maintenance.expire-snapshots.enabled' = 'true';
SET 'flink-maintenance.delete-orphan-files.enabled' = 'true';

-- Configure rewrite data files
SET 'flink-maintenance.rewrite.max-bytes' = '1073741824';

-- Configure expire snapshots
SET 'flink-maintenance.expire-snapshots.retain-last' = '5';

-- Configure delete orphan files
SET 'flink-maintenance.delete-orphan-files.min-age-seconds' = '259200';

-- Configure maintenance lock (JDBC)
SET 'flink-maintenance.lock.type' = 'jdbc';
SET 'flink-maintenance.lock.lock-id' = 'catalog.db.table';
SET 'flink-maintenance.lock.jdbc.uri' = 'jdbc:postgresql://localhost:5432/iceberg';
SET 'flink-maintenance.lock.jdbc.init-lock-tables' = 'true';

-- Now run writes; maintenance will be scheduled post-commit
INSERT INTO db.tbl SELECT ...;

或者在表 DDL 中指定选项:

CREATE TABLE db.tbl (
  ...
) WITH (
  'connector' = 'iceberg',
  'catalog-name' = 'my_catalog',
  'catalog-database' = 'db',
  'catalog-table' = 'tbl',
  'flink-maintenance.rewrite.enabled' = 'true',
  'flink-maintenance.expire-snapshots.enabled' = 'true',
  'flink-maintenance.delete-orphan-files.enabled' = 'true',

  'flink-maintenance.rewrite.max-bytes' = '1073741824',
  'flink-maintenance.expire-snapshots.retain-last' = '5',
  'flink-maintenance.delete-orphan-files.min-age-seconds' = '259200',

  'flink-maintenance.lock.type' = 'jdbc',
  'flink-maintenance.lock.lock-id' = 'catalog.db.table',
  'flink-maintenance.lock.jdbc.uri' = 'jdbc:postgresql://localhost:5432/iceberg',
  'flink-maintenance.lock.jdbc.init-lock-tables' = 'true'
);

IcebergSink 维护配置(SQL)

这些配置项可通过 SQL(SET 语句或表的 WITH 选项)设置,也可通过 IcebergSink.Builder.set() / setAll() 设置。

启用开关

配置项说明默认值
flink-maintenance.rewrite.enabled启用压缩(重写数据文件)false
flink-maintenance.expire-snapshots.enabled启用快照过期false
flink-maintenance.delete-orphan-files.enabled启用孤儿文件删除false

数据文件重写配置

配置项说明默认值
flink-maintenance.rewrite.schedule.commit-count累计 N 次提交后触发10
flink-maintenance.rewrite.schedule.data-file-count累计 N 个数据文件后触发1000
flink-maintenance.rewrite.schedule.data-file-size数据文件总大小(字节)达到后触发107374182400(100GB)
flink-maintenance.rewrite.schedule.interval-second达到时间间隔(秒)后触发600
flink-maintenance.rewrite.max-bytes每次执行允许重写的最大字节数Long.MAX_VALUE
flink-maintenance.rewrite.partial-progress.enabled启用部分进度提交false
flink-maintenance.rewrite.partial-progress.max-commits部分进度提交的最大次数10

快照过期配置

键描述默认值
flink-maintenance.expire-snapshots.schedule.commit-count在 N 次提交后触发10
flink-maintenance.expire-snapshots.schedule.interval-second经过时间间隔(秒)后触发3600(1 小时)
flink-maintenance.expire-snapshots.max-snapshot-age-seconds要保留的快照的最大存活时间(秒)未设置
flink-maintenance.expire-snapshots.retain-last要保留的快照的最小数量未设置
flink-maintenance.expire-snapshots.delete-batch-size删除过期文件的批处理大小1000
flink-maintenance.expire-snapshots.clean-expired-metadata移除过期元数据(分区规范、Schema)true
flink-maintenance.expire-snapshots.planning-worker-pool-size规划所用的工作线程池大小共享线程池

删除孤儿文件配置

键描述默认值
flink-maintenance.delete-orphan-files.schedule.interval-second经过时间间隔(秒)后触发3600(1 小时)
flink-maintenance.delete-orphan-files.min-age-seconds参与删除的文件的最小存活时间(秒)259200(3 天)
flink-maintenance.delete-orphan-files.delete-batch-size删除孤儿文件的批处理大小1000
flink-maintenance.delete-orphan-files.location开始递归列举的位置表位置
flink-maintenance.delete-orphan-files.use-prefix-listing使用前缀列举来发现文件true
flink-maintenance.delete-orphan-files.planning-worker-pool-size规划所用的工作线程池大小共享线程池
flink-maintenance.delete-orphan-files.equal-schemes等价的 scheme(格式:s3n=s3,s3a=s3)s3n=s3,s3a=s3
flink-maintenance.delete-orphan-files.equal-authorities等价的 authority(格式:auth1=auth2)未设置
flink-maintenance.delete-orphan-files.prefix-mismatch-mode前缀不匹配时的行为:ERROR、IGNORE、DELETEERROR

锁配置(SQL)

这些键用于 SQL(SET 或表的 WITH 选项),在启用维护进行写入时生效。

  • JDBC
键描述默认值
flink-maintenance.lock.type设置为 jdbc
flink-maintenance.lock.lock-id每张表唯一的锁 ID
flink-maintenance.lock.jdbc.uriJDBC URI
flink-maintenance.lock.jdbc.init-lock-tables自动创建锁表false
  • ZooKeeper
键描述默认值
flink-maintenance.lock.type设置为 zookeeper
flink-maintenance.lock.lock-id每张表唯一的锁 ID
flink-maintenance.lock.zookeeper.uriZooKeeper 连接 URI
flink-maintenance.lock.zookeeper.session-timeout-ms会话超时时间(毫秒)60000
flink-maintenance.lock.zookeeper.connection-timeout-ms连接超时时间(毫秒)15000
flink-maintenance.lock.zookeeper.max-retries最大重试次数3
flink-maintenance.lock.zookeeper.base-sleep-ms重试之间的基础等待时间(毫秒)3000
flink-maintenance.lock.zookeeper.max-sleep-ms重试之间的最大等待时间(毫秒)。用于限制指数退避延迟的上限。10000
flink-maintenance.lock.zookeeper.retry-policyZooKeeper 客户端的重试策略名称。支持的取值包括:ONE_TIME、N_TIME、BOUNDED_EXPONENTIAL_BACKOFF、UNTIL_ELAPSED、EXPONENTIAL_BACKOFF。EXPONENTIAL_BACKOFF
  • 协调器锁(COORDINATOR LOCK)
键描述默认值
flink-maintenance.lock.type设置为 `` 或不设置

最佳实践

资源管理

  • 为维护任务使用专用的 slot sharing group
  • 根据集群资源设置合适的并行度
  • 启用 checkpointing 以实现容错

调度策略

  • 使用 rateLimit 避免执行过于频繁
  • 对写入密集型表使用 scheduleOnCommitCount
  • 使用 scheduleOnDataFileCount 进行细粒度控制

性能调优

  • 根据存储性能调整 deleteBatchSize
  • 对大型重写操作启用 partialProgressEnabled
  • 设置合理的 maxRewriteBytes 限制
  • 设置合适的 maxFileGroupSizeBytes 可以将大型 FileGroup 拆分为更小的组,从而提升并行处理速度

故障排查

文件删除过程中出现 OutOfMemoryError

场景: 当维护任务尝试在单个批次中删除大量文件时,可能会出现此问题,尤其是在保留历史较长的表中或批量删除之后。原因: 每个文件的删除都涉及元数据和对象存储操作,这些操作共同会占用大量内存。较大的批次会放大这一影响,并可能耗尽 JVM 堆内存。建议: 减小批次大小,以限制删除过程中的内存占用。

.deleteBatchSize(500) // Example: 500 files per batch

锁冲突

场景: 在多作业或高可用环境中,两个或多个 Flink 作业可能同时尝试对同一张表执行维护操作。原因: 并发作业会竞争同一把分布式锁,从而导致重试并可能产生延迟。建议: 增大锁检查延迟和速率限制,使失败的尝试能够退避并降低争用。

.lockCheckDelay(Duration.ofMinutes(1)) // Wait longer before re-checking lock
.rateLimit(Duration.ofMinutes(10))     // Reduce frequency of task execution

慢速的重写操作

场景: 含有大量小文件的大表可能需要在单次运行中重写 TB 级别的数据,这可能会压垮可用资源。原因: 如果不加以限制,重写任务会尝试一次性处理所有符合条件的文件,从而导致执行时间过长,甚至可能造成作业失败。建议: 启用部分进度(partial progress),使重写的文件能够以较小的批次提交,并对每次执行重写的数据量设置上限。

.partialProgressEnabled(true) // Commit progress incrementally
.partialProgressMaxCommits(3) // Allow up to 3 commits per run
.maxRewriteBytes(1L * 1024 * 1024 * 1024) // Limit to ~1GB per run

评论

登录后参与评论

正在加载评论…