Flink内存分布
Flink1.10版本后为了满足更细粒度以及灵活的内存管理,升级了内存模型,对内存组成进行了比较大的调整,由于在Flink中计算主要存在于TaskManager节点,这里说的Flink内存模型也就是TaskManager的内存模型,JobManager的内存模型与TaskManager的内存模型类似。

上图是Flink内存模型,从图中可以看出Flink 进程总内存(Total Process Memory)包含了Flink总内存(Total Flink Memory)和JVM特定内存。Flink总内存又包括JVM堆内存(JVM Heap)、托管内存(Managed Memory) 、直接内存(Direct Memory)。下面分别介绍各个部分内存功能以及参数配置。
Flink堆内存(JVM Heap)
Flink堆内存就是JVM堆内存(JVM Heap),分为Framework堆内存(Framework Heap)和Task堆内存(Task Heap),其中Framework 主要用于Flink框架本身需要的内存空间,Task堆内存则用于Flink算子、用户代码执行及状态数据存储,也被称为TaskExecutor使用的内存,两者的主要区别在于是否将内存计入Slot计算资源中,Framework堆内存不会将内存分配给Slot,Task堆内存会分配给Slot。
Framework堆内存(Framework Heap)
Framework堆内存配置参数为:taskmanager.memory.framework.heap.size,该值默认为128M。
Task 堆内存(Task Heap)
Task堆内存配置参数为:taskmanager.memory.task.heap.size,该值没有默认值,如果没有指定会自动用Flink总内存减去Framework堆内存(Framework Heap)、托管内存(Managed Memory)、Framework非堆内存(Framework Off-Heap)、Task非堆内存(Task Off-Heap)、NetWork的剩余内存。
Flink非堆内存(Off-Heap Memory)
非堆内存也可以叫做堆外内存,更准确来说是大部分的堆外内存,包含了托管内存(Managed Memory)、直接内存(Direct Memory)两部分。
托管内存(Managed Memory)
托管内存(Managed Memory)是由Flink负责分配和管理的本地堆外内存,在流处理作业中用于RocksDBstateBackend状态存储后端,在批处理作业中用于排序、哈希表及缓存中间结果。
托管内存(Managed Memory)配置参数有两个,分别如下:
- taskmanager.memory.managed.fraction,默认值0.4,如果未显式指定托管内存大小,则使用总Flink内存的百分比作为托管内存。
- taskmanager.memory.managed.size,无默认值,一般也不指定,而是按照比例来推定,更加灵活。
直接内存(Direct Memory)
直接内存(Direct Memory)分为Framework非堆内存(Framework Off-Heap)、Task 非堆内存(Task Off-Heap)和Network三个部分。直接内存主要作用是减少GC压力、提升性能效率。
Framework 非堆内存(Framework Off-Heap)
Framework 非堆内存即taskexecutor的Framework 堆外内存大小,不会分配给slot,配置参数为:taskmanager.memory.framework.off-heap.size,默认值128M。
Task非堆内存(Task Off-Heap)
Task非堆内存,配置参数taskmanager.memory.task.off-heap.size,默认值为0,即不使用。
Network
Network内存存储空间主要用于基于Netty进行网络数据交换,数据传输的本地缓存,例如:TaskManager之间Shuffle、广播、与外部组件的数据传输。Network的配置相关参数有3个,分别如下:
- taskmanager.memory.network.min:网络缓存的最小值,默认64MB;
- taskmanager.memory.network.max:网络缓存的最大值,默认1GB;
- taskmanager.memory.network.fraction:网络缓存占Flink总内存taskmanager.memory.flink.size的比例,默认值0.1。若根据此比例算出的内存量比最小值小或比最大值大,就会限制到最小值或者最大值。
JVM 特定内存
JVM特定内存是JVM堆外内存的另一小部分内存,其不在Flink总内存范围之内,包括JVM元空间(JVM Metaspace)和JVM Overhead 两部分,其中JVM元空间存储JVM加载类的元数据,加载的类越多,需要的内存空间越大,该部分默认值为256M;JVM Overhead 则主要用于其他JVM开销,例如代码缓存、线程栈等,该部分默认值为TaskManager分配内存的0.1倍,最小值为192M,最大值为1GB。
Flink内存优化建议
关于Flink内存优化有如下几点建议:
- 在使用Flink过程中,我们可以设置参数来指定JobManager和TaskManager内存大小。Standalone部署模式下,可以通过 jobmanager.memory.flink.size和taskmanager.memory.flink.size 来指定JM和TM的内存大小。容器部署模式下(如:K8s,Yarn),可以通过jobmanager.memory.process.size和taskmanager.memory.process.size 来指定JM和TM的内存大小。JobManager主要负责管理TaskManager,这里给的内存可以相比TaskManager少一些,重点配置TaskManager内存。
- 提交Flink任务时指定的内存资源要保证程序不出现反压,并且在程序无反压使用内存基础上提供一些额外内存资源,这些内存资源可以在应用程序恢复期间更快的恢复状态及更快的处理停机期间累积的数据。
- Flink内存管理中,如果需要根据Flink程序做一些调整建议有限调整fraction比例参数,例如:网络缓存占比taskmanager.memory.network.fraction(根据网络流量大小调节)与托管内存占比taskmanager.memory.managed.fraction(根据RocksDB状态大小调节),这样做可以间接影响任务内存的配额,需要特别注意的是如果手动指定较多的固定参数很有可能出现内存配额冲突导致Flink程序部署失败。
指定内存提交Flink任务案例
下面以Flink Yarn Application模式提交任务为例来演示提交Flink任务时指定JobManager和TaskManager内存。
以读取socket数据统计WordCount为例,代码如下:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
/**
* Socket中数据如下:
* 001,186,187,busy,1000,10
* 002,187,186,fail,2000,20
*/
DataStreamSource<String> ds = env.socketTextStream("node5", 9999);
//对ds进行转换处理,得到StationLog对象
SingleOutputStreamOperator<StationLog> stationLogDS = ds.map(new MapFunction<String, StationLog>() {
@Override
public StationLog map(String line) throws Exception {
String[] arr = line.split(",");
return new StationLog(
arr[0].trim(),
arr[1].trim(),
arr[2].trim(),
arr[3].trim(),
Long.valueOf(arr[4]),
Long.valueOf(arr[5])
);
}
});
stationLogDS.keyBy(stationLog -> stationLog.callOut)
.map(new RichMapFunction<StationLog, String>() {
//定义ValueState 用来存放同一个主叫号码的通话总时长
private ValueState<Long> valueState;
@Override
public void open(Configuration parameters) throws Exception {
//定义状态TTL
StateTtlConfig ttlConfig = StateTtlConfig
//设置状态有效期为10秒
.newBuilder(org.apache.flink.api.common.time.Time.seconds(10))
//设置状态更新类型为OnCreateAndWrite,即状态TTL创建和写入时更新
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
//设置状态可见性为NeverReturnExpired,即状态过期后不返回
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build();
//定义状态描述器
ValueStateDescriptor<Long> descriptor = new ValueStateDescriptor<>("value-state", Long.class);
//设置状态TTL
descriptor.enableTimeToLive(ttlConfig);
//获取状态
valueState = getRuntimeContext().getState(descriptor);
}
@Override
public String map(StationLog stationLog) throws Exception {
Long stateValue = valueState.value();
if(stateValue==null){
//如果状态值为null,说明是第一次使用,直接更新状态值
valueState.update(stationLog.duration);
}else{
//如果状态值不为null,说明不是第一次使用,需要累加通话时长
valueState.update(stateValue+stationLog.duration);
}
return stationLog.callOut+"通话总时长:"+valueState.value()+"秒";
}
}).print();
env.execute();
以上代码编写完成后打包,在node5节点启动socket服务,启动Hadoop集群,以Yarn Application 模式提交Flink任务,命令如下:
./flink run-application -t yarn-application \
-p 2 \
-Dtaskmanager.numberOfTaskSlots=2 \
-Djobmanager.memory.process.size=1024mb \
-Dtaskmanager.memory.process.size=1024mb \
-c com.wubaibao.flinkjava.code.chapter12.MemoryTest \
/root/flink-jar-test/FlinkJavaCode-1.0-SNAPSHOT-jar-with-dependencies.jar
任务提交后,可以通过查看Flink WebUI观察内存分布情况如下:
