有趣生活

当前位置:首页>科技>hadoop大数据分析基础大数据Hadoop之Flink的状态管理和容错机制

hadoop大数据分析基础大数据Hadoop之Flink的状态管理和容错机制

发布时间:2026-07-25阅读(1)

导读一、Flink中的状态官方文档:WorkingwithState|ApacheFlink有状态的计算是流处理框架要实现的重要功能,因为稍复杂的流处理场景都需....一、Flink中的状态

官方文档:Working with State | Apache Flink

有状态的计算是流处理框架要实现的重要功能,因为稍复杂的流处理场景都需要记录状态,然后在新流入数据的基础上不断更新状态。下面的几个场景都需要使用流处理的状态功能:

  • 数据流中的数据有重复,想对重复数据去重,需要记录哪些数据已经流入过应用,当新数据流入时,根据已流入过的数据来判断去重。
  • 检查输入流是否符合某个特定的模式,需要将之前流入的元素以状态的形式缓存下来。比如,判断一个温度传感器数据流中的温度是否在持续上升。
  • 对一个时间窗口内的数据进行聚合分析,分析一个小时内某项指标的75分位或99分位的数值。

个状态更新和获取的流程如下图所示,一个算子子任务接收输入流,获取对应的状态,根据新的计算结果更新状态。

1)键控状态(Keyed State)

keyed state 接口提供不同类型状态的访问接口,这些状态都作用于当前输入数据的 key 下。换句话说,这些状态仅可在 KeyedStream 上使用,在Java/Scala API上可以通过 stream.keyBy(...) 得到 KeyedStream,在Python API上可以通过 stream.key_by(...) 得到 KeyedStream。

1、控件状态特点
  • 键控状态是根据输入数据流中定义的键(key)来维护和访问的
  • Flink 为每个 key 维护一个状态实例,并将具有相同键的所有数据,都分区到同一个算子任务中,这个任务会维护和处理这个 key 对应的状态
  • 当任务处理一条数据时,它会自动将状态的访问范围限定为当前数据的 key
2、键控状态类型
TAGS标签:  hadoop  数据分析  基础  数据  hadoop大数据分

Copyright © 2024 有趣生活 All Rights Reserve吉ICP备19000289号-5 TXT地图HTML地图XML地图