【发布时间】:2020-06-16 14:30:27
【问题描述】:
- 我正在使用 Flink 版本 1.10.1 和 RocksDB 后端。
- 我知道rocksdb 使用“托管内存”中的内存,但我没有为托管内存设置任何特定值。由 Flink 完成。
- 当我监控我的应用程序时,任务管理器的可用内存总是在减少(我的意思是通过
free -h测量的操作系统的可用内存)。我怀疑原因可能是 Rocksdb。 - Question_1 => 如果
ValueState的值已过期,那么rocksdb 将从其内存中删除并从本地存储目录中删除? (我的存储空间也有限) - Question_2 =>
stream.keyBy(ipAddress),如果这个ipAddress将被rocksdb 持有(我说的是keyBy 本身而不是状态),它总是放在托管内存中吗?如果不是,那么flink heap memory会增加吗?
这是我的应用程序的一般结构:
streamA = source.filter(..);
streamA2 = source2.filter(..);
streamB = streamA.keyBy(ipAddr).window().process(); // contains value state
streamC = streamA.keyBy(ipAddr).flatMap(..); // contains value state
streamD = streamA2.keyBy(ipAddr).window.process(); // contains value state
streamE = streamA.union(streamA2).keyBy(ipAddr)....
这是我的应用程序中的状态示例:
private transient ValueState<SampleObject> sampleState;
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.minutes(10))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build();
ValueStateDescriptor<SampleObject> sampleValueStateDescriptor = new ValueStateDescriptor<>(
"sampleState",
TypeInformation.of(SampleObject.class)
);
sampleValueStateDescriptor.enableTimeToLive(ttlConfig);
Rocksdb 配置:
state.backend: rocksdb
state.backend.rocksdb.checkpoint.transfer.thread.num: 6
state.backend.rocksdb.localdir: /pathTo/checkpoint_only_local
我为什么使用 Rocksdb
- 我正在使用rocksdb,因为我有一个巨大的密钥大小(想想它的IP 地址),HeapState 后端或其他无法处理。
- 我的应用程序使用rocksdb,因为我在用户定义的keyedprocess 函数中有一堆状态供将来决策。 (每个状态都有`StateTtlConfig)
注意
- 我的应用程序不需要增量检查点或任何有关保存点的东西。我不关心保存我的应用程序的所有快照。
【问题讨论】:
标签: apache-flink flink-streaming rocksdb