Flink的Keyed Process Function中能否定义多个状态描述符?
问题解答
Flink完全支持该操作,这是Flink键控状态的标准使用场景之一。
- 你只需要在
open()方法中分别为两类状态定义独立的状态描述符,分别获取对应的ValueState和MapState实例即可,两类状态彼此独立存储,RocksDB状态后端会自动为不同描述符对应的状态创建独立的列族存储数据,不会互相干扰。 - 运行时,同一个key下的两类状态会绑定到相同的键组,你可以在
processElement方法中直接对两个状态实例进行读写,不需要做额外的适配处理。
注意:不同状态的描述符名称必须唯一,不能重名,否则会出现状态类型冲突的报错。
以下是简单的示例代码片段:
// 先声明两个状态变量 private ValueState<String> valueState; private MapState<String, Long> mapState; @Override public void open(Configuration parameters) throws Exception { // 初始化ValueState ValueStateDescriptor<String> valueDesc = new ValueStateDescriptor<>("state1", String.class); valueState = getRuntimeContext().getState(valueDesc); // 初始化MapState MapStateDescriptor<String, Long> mapDesc = new MapStateDescriptor<>("state2", Types.STRING, Types.LONG); mapState = getRuntimeContext().getMapState(mapDesc); }
内容的提问来源于stack exchange,提问作者sparkless
相关产品推荐
相关产品推荐

