Flink 调优指南
全局配置
使用 Flink 时,可以在 $FLINK_HOME/conf/flink-conf.yaml 中设置一些全局配置。
并行度
| 选项名称 | 默认值 | 类型 | 说明 |
|---|---|---|---|
taskmanager.numberOfTaskSlots | 1 | Integer | 单个 TaskManager 可以运行的并行算子或用户函数实例的数量。建议将该值设置为 > 4,实际值需要根据数据量来设置 |
parallelism.default | 1 | Integer | 当任何位置都未指定并行度时使用的默认并行度(默认值:1)。例如,如果未设置 write.bucket_assign.tasks 的值,则将使用该值 |
内存
| 选项名称 | 默认值 | 类型 | 说明 |
|---|---|---|---|
jobmanager.memory.process.size | (none) | MemorySize | JobManager 的进程总内存大小。这包括 JobManager JVM 进程消耗的所有内存,由 Flink 总内存、JVM 元空间(JVM Metaspace)和 JVM 开销(JVM Overhead)组成 |
taskmanager.memory.task.heap.size | (none) | MemorySize | TaskExecutor 的任务堆内存大小。这是为写入缓存保留的 JVM 堆内存大小 |
taskmanager.memory.managed.size | (none) | MemorySize | TaskExecutor 的托管内存大小。这是由内存管理器管理的堆外内存大小,保留用于排序和 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) | String | RocksDB 在(TaskManager 上的)本地存放其文件的目录 |
state.checkpoints.dir | (none) | String | 用于在 Flink 支持的文件系统中存储 checkpoint 数据文件和元数据的默认目录。存储路径必须能从所有参与的进程/节点(即所有 TaskManager 和 JobManager)访问,例如 hdfs 和 oss 路径 |
state.backend.incremental | false | Boolean | 该选项决定状态后端在可能的情况下是否创建增量 checkpoint。对于增量 checkpoint,只存储与上一次 checkpoint 的差异,而不是完整的 checkpoint 状态。如果状态存储设置为 rocksdb,建议开启此选项 |
表选项(Table Options)
Flink SQL 作业可以通过 WITH 子句中的选项进行配置。实际的数据源级配置如下所列。
内存(Memory)
:::note 注意
在优化内存时,我们首先需要关注内存配置、TaskManager 的数量以及写入任务的并行度(write.tasks : 4)。在确认每个写入任务都能分配到足够的内存之后,我们再尝试设置这些内存选项。
:::
| 选项名称 | 说明 | 默认值 | 备注 |
|---|---|---|---|
write.task.max.size | 写入任务的最大内存(单位 MB),当达到该阈值时,会刷新数据量最大的桶以避免 OOM。默认 1024MB | 1024D | 写入缓冲区预留的内存为 write.task.max.size - compaction.max_memory。当所有写入任务的缓冲区总量达到阈值时,内存中最大的缓冲区会被刷新 |
write.batch.size | 为了提高写入效率,Flink 写入任务会按写入桶缓存数据,直到内存达到阈值。达到阈值后,数据缓冲区将被刷新。默认 64MB | 64D | 建议使用默认设置 |
write.log_block.size | Hudi 的日志写入器在接收到数据后不会立即刷新。写入器以 LogBlock 为单位将数据刷写到磁盘。在 LogBlock 达到阈值之前,记录会以序列化字节的形式缓存在写入器中。默认 128MB | 128 | 建议使用默认设置 |
write.merge.max_memory | 如果写入类型为 COPY_ON_WRITE,Hudi 会合并增量数据与基础文件数据。增量数据会被缓存并溢写到磁盘。该阈值控制可用的最大堆内存。默认 100MB | 100 | 建议使用默认设置 |
compaction.max_memory | 与 write.merge.max_memory 相同,但发生在 compaction 期间。默认 100MB | 100 | 如果是在线 compaction,在资源充足时可以调大,例如设置为 1024MB |
并行度
| 选项名称 | 描述 | 默认值 | 备注 |
|---|---|---|---|
write.tasks | 写入任务的并行度。每个写任务按顺序写入 1 到 N 个 bucket。默认 4 | 4 | 提高并行度对小文件的数量没有影响 |
write.bucket_assign.tasks | 分配 bucket 的算子的并行度。无默认值,使用 Flink 的 parallelism.default | parallelism.default | 提高并行度也会增加 bucket 的数量,从而增加小文件(小 bucket)的数量 |
write.index_boostrap.tasks | 索引 bootstrap 的并行度。提高并行度可以加快 bootstrap 阶段的效率。bootstrap 阶段会阻塞 checkpoint,因此需要设置更多的 checkpoint 失败容忍次数。默认使用 Flink 的 parallelism.default | parallelism.default | 仅当 index.bootsrap.enabled 为 true 时生效 |
read.tasks | 读取算子(批式与流式)的并行度。默认 4 | 4 | |
compaction.tasks | 在线 compaction 的并行度。默认 4 | 4 | Online 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),默认 100MB | 100 | 如果资源充足,建议调整为 1024MB |
compaction.target_io | 每次合并的目标 IO 量(读写合计),默认 500GB | 512000 | -- |
内存优化
MOR
- 将 Flink 状态后端设置为
rocksdb(默认的in memory状态后端非常消耗内存)。 - 如果内存充足,可以将
compaction.max_memory设置得更大(默认为100MB,可调整为1024MB)。 - 注意 taskManager 分配给每个写入任务的内存,确保每个写入任务能够获得所需的内存大小
write.task.max.size。例如,taskManager 拥有4GB内存并运行两个 streamWriteFunction,则每个写入任务可分配2GB内存。请预留一定的缓冲空间,因为网络缓冲区以及 taskManager 上的其他类型任务(如 bucketAssignFunction)也会消耗内存。 - 注意 compaction 的内存变化。
compaction.max_memory控制 compaction 任务读取日志时每个任务可使用的最大内存,compaction.tasks控制 compaction 任务的并行度。
COW
- 将 Flink 状态后端设置为
rocksdb(默认的in memory状态后端非常消耗内存)。 - 同时增大
write.task.max.size和write.merge.max_memory(默认分别为1024MB和100MB,可调整为2014MB和1024MB)。 - 注意 taskManager 分配给每个写入任务的内存,确保每个写入任务能够获得所需的内存大小
write.task.max.size。例如,taskManager 拥有4GB内存并运行两个写入任务,则每个写入任务可分配2GB内存。请预留一定的缓冲空间,因为网络缓冲区以及 taskManager 上的其他类型任务(如BucketAssignFunction)也会消耗内存。
写入速率限制
在现有的数据同步方案中,快照数据 和 增量数据 会先发送到 Kafka,再由 Flink 以流式方式写入 Hudi。由于直接消费 快照数据 会导致吞吐量过高、乱序严重(随机写入分区)等问题,进而造成写入性能下降和吞吐量抖动。此时可以开启 write.rate.limit 选项,以保证写入过程平稳。
选项
| 选项名称 | 是否必填 | 默认值 | 备注 |
|---|---|---|---|
write.rate.limit | false | 0 | 默认关闭 |
受管理内存的写入缓冲区
默认情况下,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.size | Disruptor 环形缓冲区的大小(必须是 2 的幂) | 16384 | 较大的值可以吸收写入突发流量,但会占用更多堆内存 |
write.buffer.disruptor.wait.strategy | Disruptor 消费者的等待策略:BLOCKING_WAIT(默认)、SLEEPING_WAIT、YIELDING_WAIT、BUSY_SPIN_WAIT | BLOCKING_WAIT | BLOCKING_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 将其下发读取之间经过的时间(毫秒) |
sourceReaderIdleTime | source reader 的空闲时间(毫秒),即没有被分配新 split 的时间 |
这些指标通过 Flink 的标准指标系统暴露,可以转发到 Prometheus、JMX 或其他 reporter。
评论
登录后参与评论
KnowForge