运维 Hudi

Flink 调优指南

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

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

全局配置

使用 Flink 时,可以在 $FLINK_HOME/conf/flink-conf.yaml 中设置一些全局配置。

并行度

选项名称默认值类型说明
taskmanager.numberOfTaskSlots1Integer单个 TaskManager 可以运行的并行算子或用户函数实例的数量。建议将该值设置为 > 4,实际值需要根据数据量来设置
parallelism.default1Integer当任何位置都未指定并行度时使用的默认并行度(默认值:1)。例如,如果未设置 write.bucket_assign.tasks 的值,则将使用该值

内存

选项名称默认值类型说明
jobmanager.memory.process.size(none)MemorySizeJobManager 的进程总内存大小。这包括 JobManager JVM 进程消耗的所有内存,由 Flink 总内存、JVM 元空间(JVM Metaspace)和 JVM 开销(JVM Overhead)组成
taskmanager.memory.task.heap.size(none)MemorySizeTaskExecutor 的任务堆内存大小。这是为写入缓存保留的 JVM 堆内存大小
taskmanager.memory.managed.size(none)MemorySizeTaskExecutor 的托管内存大小。这是由内存管理器管理的堆外内存大小,保留用于排序和 RocksDB 状态后端。如果选择 RocksDB 作为状态后端,则需要设置该内存

检查点

选项名称默认值类型说明
execution.checkpointing.interval(none)Duration将该值设置为 execution.checkpointing.interval = 150000ms,150000ms = 2.5 分钟。配置此参数等同于启用 checkpoint
state.backend(none)String用于存储状态的状态后端。我们建议将状态存储设置为 rocksdb:state.backend: rocksdb
state.backend.rocksdb.localdir(none)StringRocksDB 在(TaskManager 上的)本地存放其文件的目录
state.checkpoints.dir(none)String用于在 Flink 支持的文件系统中存储 checkpoint 数据文件和元数据的默认目录。存储路径必须能从所有参与的进程/节点(即所有 TaskManager 和 JobManager)访问,例如 hdfs 和 oss 路径
state.backend.incrementalfalseBoolean该选项决定状态后端在可能的情况下是否创建增量 checkpoint。对于增量 checkpoint,只存储与上一次 checkpoint 的差异,而不是完整的 checkpoint 状态。如果状态存储设置为 rocksdb,建议开启此选项

表选项(Table Options)

Flink SQL 作业可以通过 WITH 子句中的选项进行配置。实际的数据源级配置如下所列。

内存(Memory)

:::note 注意
在优化内存时,我们首先需要关注内存配置、TaskManager 的数量以及写入任务的并行度(write.tasks : 4)。在确认每个写入任务都能分配到足够的内存之后,我们再尝试设置这些内存选项。
:::

选项名称说明默认值备注
write.task.max.size写入任务的最大内存(单位 MB),当达到该阈值时,会刷新数据量最大的桶以避免 OOM。默认 1024MB1024D写入缓冲区预留的内存为 write.task.max.size - compaction.max_memory。当所有写入任务的缓冲区总量达到阈值时,内存中最大的缓冲区会被刷新
write.batch.size为了提高写入效率,Flink 写入任务会按写入桶缓存数据,直到内存达到阈值。达到阈值后,数据缓冲区将被刷新。默认 64MB64D建议使用默认设置
write.log_block.sizeHudi 的日志写入器在接收到数据后不会立即刷新。写入器以 LogBlock 为单位将数据刷写到磁盘。在 LogBlock 达到阈值之前,记录会以序列化字节的形式缓存在写入器中。默认 128MB128建议使用默认设置
write.merge.max_memory如果写入类型为 COPY_ON_WRITE,Hudi 会合并增量数据与基础文件数据。增量数据会被缓存并溢写到磁盘。该阈值控制可用的最大堆内存。默认 100MB100建议使用默认设置
compaction.max_memory与 write.merge.max_memory 相同,但发生在 compaction 期间。默认 100MB100如果是在线 compaction,在资源充足时可以调大,例如设置为 1024MB

并行度

选项名称描述默认值备注
write.tasks写入任务的并行度。每个写任务按顺序写入 1 到 N 个 bucket。默认 44提高并行度对小文件的数量没有影响
write.bucket_assign.tasks分配 bucket 的算子的并行度。无默认值,使用 Flink 的 parallelism.defaultparallelism.default提高并行度也会增加 bucket 的数量,从而增加小文件(小 bucket)的数量
write.index_boostrap.tasks索引 bootstrap 的并行度。提高并行度可以加快 bootstrap 阶段的效率。bootstrap 阶段会阻塞 checkpoint,因此需要设置更多的 checkpoint 失败容忍次数。默认使用 Flink 的 parallelism.defaultparallelism.default仅当 index.bootsrap.enabled 为 true 时生效
read.tasks读取算子(批式与流式)的并行度。默认 44
compaction.tasks在线 compaction 的并行度。默认 44Online compaction 会占用写入任务的资源。建议使用 offline compaction

Compaction

note

以下选项仅适用于 online compaction(在线压缩)。

note

通过设置 compaction.async.enabled = false 可以关闭在线压缩,但我们仍然建议为写入作业开启 compaction.schedule.enable。之后你可以通过 offline compaction(离线压缩)来执行压缩计划。

选项名称说明默认值备注
compaction.schedule.enabled是否定期生成合并计划true建议开启,即使 compaction.async.enabled = false
compaction.async.enabled异步合并,MOR 表默认启用true关闭该选项即可关闭 在线合并(online compaction)
compaction.trigger.strategy触发合并的策略num_commits可选值:num_commits:当达到 N 次增量提交时触发合并;time_elapsed:当距离上次合并已过去的时间超过 N 秒时触发合并;num_and_time:当同时满足 NUM_COMMITS 和 TIME_ELAPSED 时触发合并;num_or_time:当满足 NUM_COMMITS 或 TIME_ELAPSED 之一时触发合并。
compaction.delta_commits触发合并所需的最大增量提交次数,默认 5 次提交5--
compaction.delta_seconds触发合并所需的最大增量时间,默认 1 小时3600--
compaction.max_memory合并可溢写映射(spillable map)的最大内存(单位 MB),默认 100MB100如果资源充足,建议调整为 1024MB
compaction.target_io每次合并的目标 IO 量(读写合计),默认 500GB512000--

内存优化

MOR

  1. 将 Flink 状态后端设置为 rocksdb(默认的 in memory 状态后端非常消耗内存)。
  2. 如果内存充足,可以将 compaction.max_memory 设置得更大(默认为 100MB,可调整为 1024MB)。
  3. 注意 taskManager 分配给每个写入任务的内存,确保每个写入任务能够获得所需的内存大小 write.task.max.size。例如,taskManager 拥有 4GB 内存并运行两个 streamWriteFunction,则每个写入任务可分配 2GB 内存。请预留一定的缓冲空间,因为网络缓冲区以及 taskManager 上的其他类型任务(如 bucketAssignFunction)也会消耗内存。
  4. 注意 compaction 的内存变化。compaction.max_memory 控制 compaction 任务读取日志时每个任务可使用的最大内存,compaction.tasks 控制 compaction 任务的并行度。

COW

  1. 将 Flink 状态后端设置为 rocksdb(默认的 in memory 状态后端非常消耗内存)。
  2. 同时增大 write.task.max.size 和 write.merge.max_memory(默认分别为 1024MB 和 100MB,可调整为 2014MB 和 1024MB)。
  3. 注意 taskManager 分配给每个写入任务的内存,确保每个写入任务能够获得所需的内存大小 write.task.max.size。例如,taskManager 拥有 4GB 内存并运行两个写入任务,则每个写入任务可分配 2GB 内存。请预留一定的缓冲空间,因为网络缓冲区以及 taskManager 上的其他类型任务(如 BucketAssignFunction)也会消耗内存。

写入速率限制

在现有的数据同步方案中,快照数据 和 增量数据 会先发送到 Kafka,再由 Flink 以流式方式写入 Hudi。由于直接消费 快照数据 会导致吞吐量过高、乱序严重(随机写入分区)等问题,进而造成写入性能下降和吞吐量抖动。此时可以开启 write.rate.limit 选项,以保证写入过程平稳。

选项

选项名称是否必填默认值备注
write.rate.limitfalse0默认关闭

受管理内存的写入缓冲区

默认情况下,Flink 写入缓冲区使用 JVM 堆内存(ON_HEAP)。在堆内存配额受限的容器化环境中,可以切换到 Flink 的受管理(堆外)内存池,以降低 GC 压力并避免 OOM 错误。

注意

使用 MANAGED 内存类型时,请确保在 flink-conf.yaml 中配置了足够的 taskmanager.memory.managed.size。

选项名称说明默认值备注
write.buffer.memory.type写缓冲区的内存类型:ON_HEAP(默认,使用 JVM 堆内存)或 MANAGED(使用 Flink 托管的堆外内存)ON_HEAP切换为 MANAGED,可在内存受限的部署环境中避免 OOM
write.memory.segment.page.size写缓冲区所用内存段的页大小(字节)32768(32 KB)根据工作负载特征调整;更大的页可减少大记录的开销

Disruptor 缓冲区调优

当表选项中设置 write.buffer.type=DISRUPTOR 时(参见 Append Write Buffer),以下调优选项用于控制 Disruptor 环形缓冲区:

选项名称说明默认值备注
write.buffer.disruptor.ring.sizeDisruptor 环形缓冲区的大小(必须是 2 的幂)16384较大的值可以吸收写入突发流量,但会占用更多堆内存
write.buffer.disruptor.wait.strategyDisruptor 消费者的等待策略:BLOCKING_WAIT(默认)、SLEEPING_WAIT、YIELDING_WAIT、BUSY_SPIN_WAITBLOCKING_WAITBLOCKING_WAIT 在容器化环境中最安全;BUSY_SPIN_WAIT 延迟最低,但需要占用一个专用 CPU 核心

基于时间线服务的 Marker

从 Hudi 1.2.0 起,Flink 写入器支持 TIMELINE_SERVER_BASED 类型的 marker(hoodie.write.markers.type=TIMELINE_SERVER_BASED)。在对象存储(S3、GCS、ADLS)上,推荐使用该方式而非 DIRECT marker,因为目录列表操作成本高昂,会使 DIRECT marker 变慢。

CREATE TABLE my_table (...)
WITH (
  'connector' = 'hudi',
  'path' = 's3a://my-bucket/my-table',
  'hoodie.write.markers.type' = 'TIMELINE_SERVER_BASED'
  -- other options
);

Source V2 读取延迟指标

当启用 Source V2(read.source-v2.enabled=true)时,系统会输出以下读取延迟指标,以便监控流式管道的健康状况:

指标说明
issuedInstantDelay从写入新 instant 到 source 将其下发读取之间经过的时间(毫秒)
sourceReaderIdleTimesource reader 的空闲时间(毫秒),即没有被分配新 split 的时间

这些指标通过 Flink 的标准指标系统暴露,可以转发到 Prometheus、JMX 或其他 reporter。

评论

登录后参与评论

正在加载评论…