【发布时间】:2021-04-23 16:20:42
【问题描述】:
我有以下flink keyedprocess函数。我基本上是在尝试实现状态设计模式。
public AlertProcessor extends KeyedProcessFunction<Tuple2<String, String>, Event1, Event2> {
private transient AlertState currentState;
private transient AlertState activeAlertState;
private transient AlertState noActiveAlertState;
private transient AlertState resolvedAlertState;
@Override
public void open(Configuration parameters) {
activeAlertState = new ActiveAlertState();
noActiveAlertState = new NoActiveAlertState();
resolvedAlertState = new ResolvedAlertState();
}
@Override
public processElement(Event1 event1, Context ctx, Collector<Event2> out) throws Exception {
// Would the below if condition work for multiple keys?
if (currentAlertState == null) {
currentAlertState = noActiveAlertState;
}
currentAlertState.handle(event1, out);
}
private interface AlertState {
void handle(Event1 event1, Collector<Event2> out);
}
private class ActiveAlertState implements AlertState {
void handle(Event1 event1, Collector<Event2> out) {
logger.debug("Moving to no alertState");
// Do something and push some Event2 to out
currentAlertState = resolvedActiveAlertState;
}
}
private class NoActiveAlertState implements AlertState {
void handle(Event1 event1, Collector<Event2> out) {
logger.debug("Moving to no alertState");
// Do something and push some Event2 to out
currentAlertState = activeAlertState;
}
}
private class ResolvedAlertState implements AlertState {
void handle(Event1 event1, Collector<Event2> out) {
logger.debug("Moving to no alertState");
// Do something and push some Event2 to out
currentAlertState = noActiveAlertState;
}
}
}
我的问题是-
- 流中的每个键是否会有一个 AlertProcessor 实例(或对象)?换句话说, currentAlertState 对象是否每个键都是唯一的?或者这个 AlertProcessor 操作符的每个实例都会有一个 currentAlertState?
如果 currentAlertState 是运算符的每个实例,那么此代码将不会真正起作用,因为 currentAlertState 将被不同的键覆盖。我的理解正确吗?
-
我可以将 currentAlertState 存储为键控状态,并为每个 processElement() 调用初始化它。如果这样做,我不需要在 handle() 实现中将 currentAlertState 分配或设置为下一个状态,因为 currentAlertState 无论如何都会根据 flink 状态进行初始化。
-
有没有更好的方法在 flink 中实现状态设计模式并且仍然减少创建的状态对象的数量?
【问题讨论】:
标签: apache-flink flink-streaming