【问题标题】:Flink ValueState will be removed from storage after expired when using Rocksdb?使用 Rocksdb 时 Flink ValueState 过期后会从存储中移除?
【发布时间】: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


    【解决方案1】:

    Flink ValueState 使用时过期后会从存储中移除 Rocksdb?

    是的,但不是立即。 (而在 Flink 的一些早期版本中,答案是“视情况而定”。)

    在您的状态 ttl 配置中,您没有指定您希望如何完成状态清理。在这种情况下,过期值会在读取时显式删除(例如ValueState#value),否则会在后台定期进行垃圾收集。在 RocksDB 的情况下,这个后台清理是在压缩期间完成的。换句话说,清理不是立即的。 docs 提供了有关如何调整它的更多详细信息 - 您可以将清理配置为更快完成,但会降低一些性能。

    keyBy 本身不使用任何状态。键选择器功能用于对流进行分区,但键不与 keyBy 关联存储。只有窗口和平面图操作保持状态,这是每个键的状态,所有这些键状态都将在 RocksDB 中(除非您已将定时器配置为堆上,这是一个选项,但在 Flink 1.10 定时器中默认情况下存储在堆外,在rocksdb中)。

    您可以将 flatmap 更改为 KeyedProcessFunction 并使用计时器显式清除状态键的状态 - 这将使您可以直接控制清除状态的确切时间,而不是依赖状态 TTL 机制最终清除状态。

    但更有可能的是,窗户正在建立相当大的状态。如果您可以切换到进行预聚合(通过reduceaggregate)可能会有很大帮助。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多