【发布时间】:2020-12-11 16:29:07
【问题描述】:
我想通过processKeyedFunction来实现aggregationFunction,因为默认aggregationFunction不支持丰富的功能, 另外,我尝试了aggreagationFunction + processWindowFunction(https://ci.apache.org/projects/flink/flink-docs-stable/dev/stream/operators/windows.html),但它也不能满足我的需求,所以我必须使用基本的processKeyedFunction来实现aggregationFunction,我的问题细节如下:
在processFunction中,我定义了一个windowState用于stage元素的聚合值,代码如下:
public void open(Configuration parameters) throws Exception {
followCacheMap = FollowSet.getInstance();
windowState = getRuntimeContext().getMapState(windowStateDescriptor);
currentTimer = getRuntimeContext().getState(new ValueStateDescriptor<Long>(
"timer",
Long.class
));
在processElement()函数中,我使用windowState(这是一个MapState在open函数中初始化)聚合窗口元素,并注册第一次timeServie清除当前窗口状态,代码如下:
@Override
public void processElement(FollowData value, Context ctx, Collector<FollowData> out) throws Exception
{
if ( (currentTimer==null || (currentTimer.value() ==null) || (long)currentTimer.value()==0 ) && value.getClickTime() != null) {
currentTimer.update(value.getClickTime() + interval);
ctx.timerService().registerEventTimeTimer((long)currentTimer.value());
}
windowState = doMyAggregation(value);
}
在onTimer()函数中,首先我在next 1分钟内注册下一次timeService,并清除窗口状态
@Override
public void onTimer(long timestamp, OnTimerContext ctx, Collector<FollowData> out) throws Exception {
currentTimer.update(timestamp + interval); // interval is 1 minute
ctx.timerService().registerEventTimeTimer((long)currentTimer.value());
out.collect(windowState);
windowState.clear();
}
但是程序运行的时候发现onTimer里面的windowState都是空的,但是processElement()函数里面却不是空的,不知道为什么会这样,可能是执行逻辑不一样,怎么可能我解决这个, 提前致谢!
新增关于 doMyAggregation() 部分的代码
windowState是一个MapState,key是“mykey”,value是一个自定义的Object AggregateFollow
public class AggregateFollow {
private String clicked;
private String unionid;
private ArrayList allFollows;
private int enterCnt;
private Long clickTime;
}
和doMyAggregation(value)函数差不多就是这样,doMyAggregation的作用是获取源字段为'follow'的所有值,但如果1分钟内没有字段为'click'的值, 'follow'值应该是过时的,一句话,就像'follow'数据和'click'数据的join操作,
AggregateFollow acc = windowState.get(windowkey);
String flag = acc.getClicked();
ArrayList<FollowData> followDataList = acc.getAllFollows();
if ("0".equals(flag)) {
if ("follow".equals(value.getSource())) {
followDataList.add(value);
acc.setAllFollows(followDataList);
}
if ("click".equals(value.getSource())) {
String unionid = value.getUnionid();
clickTime = value.getClickTime();
if (followDataList.size() > 0) {
ArrayList listNew = new ArrayList();
for (FollowData followData : followDataList) {
followData.setUnionid(unionid);
followData.setClickTime(clickTime);
followData.setSource("joined_flag"); //
}
acc.setAllFollows(listNew);
}
acc.setClicked("1");
acc.setUnionid(unionid);
acc.setClickTime(clickTime);
windowState.put(windowkey, acc);
}
} else if ("1".equals(flag)) {
if ("follow".equals(value.getSource())) {
value.setUnionid(acc.getUnionid());
value.setClickTime(acc.getClickTime());
value.setSource("joined_flag");
followDataList.add(value);
acc.setAllFollows(followDataList);
windowState.put(windowkey, acc);
}
}
由于性能问题,原来的 windowAPI 对我来说不是一个有效的选择,我认为这里唯一的方法是使用 processFunction + ontimer 和 Guava Cache , 非常感谢
【问题讨论】:
标签: apache-flink