欢迎光临
我们一直在努力

Flink从入门到上天系列第十七篇:Flink当中的算子状态

一:算子状态

        算子状态:一个算子,会有多个并行子任务。作用范围被限定为当前算子任务。算子状态跟数据的key无关,所以不同key的数据只要被分发到同一个并行子任务,就会访问到同一个Operator State。

        算子状态的实际应用场景不如Keyed State多,一般用在Source或Sink等与外部系统连接的算子上,或者完全没有key定义的场景。比如Flink的Kafka连接器中,就用到了算子状态。

        当算子的并行度发生变化时,算子状态也支持在并行的算子任务实例之间做重组分配。根据状态的类型不同,重组分配的方案也会不同。

        算子状态也支持不同的结构类型,主要有三种:ListState、UnionListState和BroadcastState。

二:列表状态

1:使用普通变量实现

// 需求,在map算子上,计算每一个并行度处理多少数据
public class Flink08_OpeState_Var {
public static void main(String[] args) throws Exception {
// 1. 准备环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 2. 配置并行度
env.setParallelism(2);
// 3. 读取数据
DataStreamSource<String> socketDs = env.socketTextStream("bigdata137", 8888);
// 4. 数据类型转换
SingleOutputStreamOperator<WaterSensor> wsDs = socketDs.map(new WaterSensorMapFunction());
SingleOutputStreamOperator<String> map = wsDs.map(new RichMapFunction<WaterSensor, String>() {

Integer count = 0;

@Override
public String map(WaterSensor value) throws Exception {
// 使用复函数,可以通过复函数获取更丰富的上下文信息。
return "并行子任务"+getRuntimeContext().getIndexOfThisSubtask() + "处理了"+ ++count+"条数据";
}
});

//7. 打印输出
map.print();
//8. 提交作业
env.execute();
}
}

2:案例分析

1> 并行子任务0处理了1条数据
2> 并行子任务1处理了1条数据
1> 并行子任务0处理了2条数据
2> 并行子任务1处理了2条数据
1> 并行子任务0处理了3条数据
2> 并行子任务1处理了3条数据
1> 并行子任务0处理了4条数据
2> 并行子任务1处理了4条数据

[root@bigdata137 ~]# nc -lk 8888
ws1,1,1
ws1,1,1
ws1,1,1
ws1,1,1
ws1,1,1
ws1,1,1
ws1,1,1
ws1,1,1

基于一个integer就已经实现了需求,并且并行度之间是隔离的。

严格来讲,状态和普通变量就差别一个序列化能力。

如果三个并行度变成了两个并行度。把并行子任务上的状态数据,重新分配给扩缩之后的并行度子任务。

例如:

列表状态:上游三个算子【1,2】【3,4】【5,6】缩容两个算子【135】【246】

联合列表:上游三个算子【1,2】【3,4】【5,6】缩容两个算子【123456】【123456】

3:算子状态实现

        把你的普通变量放到状态上去,用状态的能力给你持久化,和做容灾回复。

package com.dashu.day08;

import com.dashu.bean.WaterSensor;
import com.dashu.bean.WaterSensorMapFunction;
import org.apache.commons.compress.archivers.dump.DumpArchiveEntry;
import org.apache.flink.api.common.functions.RichMapFunction;
import org.apache.flink.api.common.state.ListState;
import org.apache.flink.api.common.state.ListStateDescriptor;
import org.apache.flink.api.common.state.OperatorStateStore;
import org.apache.flink.runtime.state.FunctionInitializationContext;
import org.apache.flink.runtime.state.FunctionSnapshotContext;
import org.apache.flink.streaming.api.checkpoint.CheckpointedFunction;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

// 需求,在map算子上,计算每一个并行度处理多少数据
public class Flink09_OpeState_State {
public static void main(String[] args) throws Exception {
// 1. 准备环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 2. 配置并行度
env.setParallelism(2);
env.enableCheckpointing(10000);
// 3. 读取数据
DataStreamSource<String> socketDs = env.socketTextStream("bigdata137", 8888);
// 4. 数据类型转换
SingleOutputStreamOperator<WaterSensor> wsDs = socketDs.map(new WaterSensorMapFunction());
SingleOutputStreamOperator<String> map = wsDs.map(new MyMap());

//7. 打印输出
map.print();
//8. 提交作业
env.execute();
}
}

class MyMap extends RichMapFunction<WaterSensor,String> implements CheckpointedFunction {

Integer count = 0;//每个并行度中都有。
ListState<Integer> countState;//每个并行度中都有

@Override
public String map(WaterSensor value) throws Exception {
return "并行子任务"+getRuntimeContext().getIndexOfThisSubtask() + "处理了"+ ++count+"条数据";
}

@Override
// 开启检查点,检查点执行的时候,会调用这个方法。
public void snapshotState(FunctionSnapshotContext context) throws Exception {
System.out.println(" ================= snapshotState ================= ");
//放之前,先给他清空
countState.clear();
//
countState.add(count);

}

@Override
//初始化状态
public void initializeState(FunctionInitializationContext context) throws Exception {
System.out.println(" ================= initializeState ================= ");
OperatorStateStore operatorStateStore = context.getOperatorStateStore();
// 获取列表状态
ListStateDescriptor listStateDescriptor = new ListStateDescriptor("countState", Integer.class);
countState= operatorStateStore.getListState(listStateDescriptor);

// 容错机制。数据恢复。
if(context.isRestored()){
Integer next = countState.get().iterator().next();
count = next;
}
}
// ================= initializeState =================
// ================= initializeState =================
// ================= snapshotState =================
// ================= snapshotState =================
//2> 并行子任务1处理了1条数据
//1> 并行子任务0处理了1条数据
// ================= snapshotState =================
// ================= snapshotState =================
//2> 并行子任务1处理了2条数据
// ================= snapshotState =================
// ================= snapshotState =================
// ================= snapshotState =================
// ================= snapshotState =================
// ================= snapshotState =================
// ================= snapshotState =================
}

一般什么时候,使用算子状态,从source当中读取数据的时候,或者sink到其他地方的时候的会用到状态。

所以flink中的kafkasource里边维护了状态。也及时kafkasoure的底层,维护了状态,保障了偏移量。

三:联合列表状态

1:两个列表状态的区别

先明确 UnionListState 和你原有 ListState 的核心区别

状态类型并行度调整时的行为适用场景
ListState 所有状态项合并后轮询均分给新并行子任务 统计、分片计算(需均分状态)
UnionListState 所有状态项合并后广播给每个新并行子任务 全量数据依赖(每个子任务都需要完整状态)

2:案例分析

/**
* UnionListState 案例:
* 1. 每个并行子任务记录自己处理过的所有 WaterSensor 的 id
* 2. 并行度调整/作业重启后,每个新并行子任务能拿到所有历史 id 列表(广播特性)
*/
public class Flink10_OperState_UnionListState {
public static void main(String[] args) throws Exception {
// 1. 准备环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 2. 配置并行度(可以先设2,测试后调整为3验证Union特性)
env.setParallelism(2);
// 开启Checkpoint,保证状态持久化(每10秒一次)
env.enableCheckpointing(10000);

// 3. 读取socket数据
DataStreamSource<String> socketDs = env.socketTextStream("bigdata137", 8888);

// 4. 转换为WaterSensor类型
SingleOutputStreamOperator<WaterSensor> wsDs = socketDs.map(new WaterSensorMapFunction());

// 5. 自定义Map算子,使用UnionListState记录所有id
SingleOutputStreamOperator<String> resultDs = wsDs.map(new MyUnionMap());

// 6. 打印输出
resultDs.print();

// 7. 提交作业
env.execute("UnionListState Demo");
}

/**
* 自定义RichMapFunction,实现CheckpointedFunction
* 核心:使用UnionListState存储所有处理过的WaterSensor id
*/
static class MyUnionMap extends RichMapFunction<WaterSensor, String> implements CheckpointedFunction {

// 存储当前并行子任务处理过的所有id(内存临时存储)
private List<String> currentIdList;
// UnionListState:算子状态,并行度调整时广播所有状态项
private UnionListState<String> unionIdState;

/**
* 初始化:创建内存列表,初始化UnionListState
*/
@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
System.out.println("===== 并行子任务" + getRuntimeContext().getIndexOfThisSubtask() + " 初始化状态 =====");

// 1. 初始化内存临时列表
currentIdList = new ArrayList<>();

// 2. 获取算子状态存储
OperatorStateStore operatorStateStore = context.getOperatorStateStore();

// 3. 定义UnionListState描述符
ListStateDescriptor<String> stateDesc = new ListStateDescriptor<>(
"union_water_sensor_id", // 状态名称
String.class // 状态数据类型
);

// 4. 获取UnionListState(关键:getUnionListState 而非 getListState)
unionIdState = operatorStateStore.getUnionListState(stateDesc);

// 5. 容错恢复:作业重启时,从UnionListState中读取所有id到内存列表
if (context.isRestored()) {
System.out.println("===== 并行子任务" + getRuntimeContext().getIndexOfThisSubtask() + " 恢复状态 =====");
// 遍历UnionListState,将所有历史id加入内存列表
Iterator<String> stateIterator = unionIdState.get().iterator();
while (stateIterator.hasNext()) {
currentIdList.add(stateIterator.next());
}
}
}

/**
* 核心处理逻辑:每条数据加入内存列表,返回当前子任务的id列表信息
*/
@Override
public String map(WaterSensor value) throws Exception {
// 1. 获取当前WaterSensor的id
String sensorId = value.getId();

// 2. 将id加入内存临时列表
currentIdList.add(sensorId);

// 3. 构造返回结果:当前子任务ID + 已处理的id数量 + 所有id列表
return "并行子任务" + getRuntimeContext().getIndexOfThisSubtask()
+ " | 累计处理id数:" + currentIdList.size()
+ " | 所有id:" + currentIdList;
}

/**
* Checkpoint触发时:将内存列表中的所有id写入UnionListState持久化
*/
@Override
public void snapshotState(FunctionSnapshotContext context) throws Exception {
System.out.println("===== 并行子任务" + getRuntimeContext().getIndexOfThisSubtask() + " 持久化状态 =====");

// 1. 清空原有状态(避免重复存储)
unionIdState.clear();

// 2. 将当前内存列表中的所有id写入UnionListState
unionIdState.addAll(currentIdList);
}
}
}

3:总结

托管状态:
描述: 由Flink框架管理状态的存储、序列化以及容错恢复等
算子状态:
作用范围: 算子的每一个并行子任务上(分区、并行度、slot)
类型:
– ListState
– UnionListState
– BroadcastState(重点)
使用步骤:
– 类必须实现CheckpointedFunction接口
– 接口方法:
initializeState: 初始化状态
snapshotState: 对状态进行备份
特点:
– 声明位置: 与普通成员变量声明位置相同
– 作用范围: 与普通成员变量作用范围相同
– 持久化: 可被持久化
– 底层原理: snapshotState阶段将普通变量放入状态中持久化,相当于普通变量也被持久保存
键控状态:
作用范围: 经过keyBy之后的每一个组,组和组之间状态是隔离的
类型:
– ValueState
– ListState
– MapState
– ReducingState
– AggregatingState
使用步骤:
– 声明位置: 在处理函数类成员变量位置声明状态
– 注意: 虽然在成员变量位置声明,但是作用范围是keyBy后的每一个组
– 初始化: 在open方法中对状态进行初始化
– 使用: 在具体的处理函数中使用状态

四:广播状态(重点)

        广播状态也是算子状态的一种,作用范围是算子子任务。他的使用方式是固定的。

        要处理的数据是一条流,他匹配的规则是另外一条流。多个流可能会有多个并行度,多个并行度需要共享规则流中数据。每一个并行度上都必须保留规则数据。规则数据直接放到广播状态里边即可,进行存储。

        我们需要处理一条流的数据,需要用到另外一条流的规则。这个时候需要用到广播状态,但是如果两个流,数据两都特别大,这个时候,就不适合使用这个状态了。

        广播状态使用非常固定:不像是其他状态那么灵活。

1:大概意思

        数据流:首先基于keyby进行分区

// 将图形使用颜色进行划分
KeyedStream<Item, Color> colorPartitionedStream = itemStream
.keyBy(new KeySelector<Item, Color>() {….});

       规则流生成广播状态:

// 一个 map descriptor,它描述了用于存储规则名称与规则本身的 map 存储结构
MapStateDescriptor<String, Rule> ruleStateDescriptor = new MapStateDescriptor<>(
"RulesBroadcastState",
BasicTypeInfo.STRING_TYPE_INFO,
TypeInformation.of(new TypeHint<Rule>() {}));

// 广播流,广播规则并且创建 broadcast state
BroadcastStream<Rule> ruleBroadcastStream = ruleStream
.broadcast(ruleStateDescriptor);

        使用规则来筛选数据,需要:关联两个流+完成模式识别过滤逻辑

        为了关联一个非广播流(keyed 或者 non-keyed)与一个广播流(BroadcastStream),我们可以调用非广播流的方法 connect(),并将 BroadcastStream 当做参数传入。这个方法的返回参数是 BroadcastConnectedStream,具有类型方法 process(),传入一个特殊的 CoProcessFunction 来书写我们的模式识别逻辑。

具体传入 process() 的是哪个类型取决于非广播流的类型:

  • 如果流是一个 keyed 流,那就对应 KeyedBroadcastProcessFunction 类型;
  • 如果流是一个 non-keyed 流,那就对应 BroadcastProcessFunction 类型。

注意:connect() 方法需要由非广播流来进行调用,BroadcastStream 作为参数传入。

        在例子中,图形流是一个 keyed stream,因此会使用 KeyedBroadcastProcessFunction 来处理。

DataStream<String> output = colorPartitionedStream
.connect(ruleBroadcastStream)
.process(

// KeyedBroadcastProcessFunction 中的类型参数表示:
// 1. key stream 中的 key 类型
// 2. 非广播流中的元素类型
// 3. 广播流中的元素类型
// 4. 结果的类型,在这里是 string

new KeyedBroadcastProcessFunction<Color, Item, Rule, String>() {
// 模式匹配逻辑
}
);

2:需求

水位超过指定阈值,发出告警,阈值可以动态修改

3:编写代码

package com.dashu.day08;

import com.dashu.bean.WaterSensor;
import com.dashu.bean.WaterSensorMapFunction;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.state.BroadcastState;
import org.apache.flink.api.common.state.MapStateDescriptor;
import org.apache.flink.api.common.state.ReadOnlyBroadcastState;
import org.apache.flink.streaming.api.datastream.BroadcastConnectedStream;
import org.apache.flink.streaming.api.datastream.BroadcastStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.co.BroadcastProcessFunction;
import org.apache.flink.util.Collector;

public class Flink10_OpenState_BrodeCat {

public static void main(String[] args) throws Exception {

// 1. 准备环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 2. 配置并行度
env.setParallelism(2);
env.enableCheckpointing(10000);
// 3. 读取数据
DataStreamSource<String> socketDs = env.socketTextStream("bigdata137", 8888);
// 4. 数据类型转换
SingleOutputStreamOperator<WaterSensor> wsDs = socketDs.map(new WaterSensorMapFunction());
//TODO 4.从指定的网络端口读取阈值信息并进行类型转换 String->Integer
SingleOutputStreamOperator<Integer> rule = env.socketTextStream("bigdata137", 8889)
.map(new MapFunction<String, Integer>() {
@Override
public Integer map(String value) throws Exception {
return Integer.valueOf(value);
}
});
//TODO 5.对阈值流数据进行广播 并声明广播状态描述器
MapStateDescriptor<String, Integer> mapStateDescriptor
= new MapStateDescriptor<String, Integer>(
"mapStateDescriptor",
String.class,
Integer.class
);
// 生成广播流
BroadcastStream<Integer> broadcast = rule.broadcast(mapStateDescriptor);
//TODO 6.关联非广播流(水位信息)以及广播流(阈值)—connect
BroadcastConnectedStream<WaterSensor, Integer> connect = wsDs.connect(broadcast);

//TODO 7.对关联后的数据进行处理—process

// 对获取到的连接流进行处理
SingleOutputStreamOperator<String> process = connect.process(new BroadcastProcessFunction<WaterSensor, Integer, String>() {

@Override
// processElement: 处理非广播流数据 从广播状态中获取阈值信息判断是否超过警戒线
public void processElement(WaterSensor value, BroadcastProcessFunction<WaterSensor, Integer, String>.ReadOnlyContext ctx, Collector<String> out) throws Exception {
// 获取广播状态
ReadOnlyBroadcastState<String, Integer> broadcastState = ctx.getBroadcastState(mapStateDescriptor);
// 从广播状态中获取水位信息
Integer i = broadcastState.get("threshHold");
i = i==null?0:i;
// 获取当前的水位值
Integer vc = value.getVc();
if (vc > i) {
out.collect("当前水位" + vc + "超过阈值" + i);
}
}

@Override
// processBroadcastElement: 处理广播流数据 将广播流中的阈值信息放到广播状态中
public void processBroadcastElement(Integer value, BroadcastProcessFunction<WaterSensor, Integer, String>.Context ctx, Collector<String> out) throws Exception {
// 获取广播状态
BroadcastState<String, Integer> broadcastState = ctx.getBroadcastState(mapStateDescriptor);
// 将广播中的阈值信息放到广播状态中
broadcastState.put("threshHold", value);
}
});

//TODO 8.打印
process.print();
//TODO 9.提交作业
env.execute();
}
}

4:案例分析

输入数据

Last login: Mon Mar 16 12:44:35 2026 from 192.168.67.1
[root@bigdata137 ~]# nc -lk 8888
ws1,1,10
ws1,1,40
ws1,1,50
ws1,1,10

Last login: Tue Mar 17 10:09:20 2026 from 192.168.67.1
[root@bigdata137 ~]# nc -lk 8889
30
8

SLF4J: Failed to load class "org.slf4j.impl.StaticLoggerBinder".
SLF4J: Defaulting to no-operation (NOP) logger implementation
SLF4J: See http://www.slf4j.org/codes.html#StaticLoggerBinder for further details.
2> 当前水位40超过阈值30
1> 当前水位50超过阈值30
2> 当前水位10超过阈值8

赞(0)
未经允许不得转载:171主机测评 » Flink从入门到上天系列第十七篇:Flink当中的算子状态
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址