【发布时间】:2020-08-27 02:18:31
【问题描述】:
我有一个键控数据流,如下所示:
{
summary:Integer
uid:String
key:String
.....
}
我需要在某个时间范围内聚合汇总值,一旦达到特定数字,将汇总和影响汇总的所有 UID 刷新到数据库/日志文件。
第一次刷新后,我想从内存中删除所有 uid,并立即刷新每个新项目。
所以我尝试了这个聚合函数。
public class AggFunc implements AggregateFunction<Item, Acc, Tuple2<Integer,List<String>>>{
private static final long serialVersionUID = 1L;
@Override
public Acc createAccumulator() {
return new Acc());
}
@Override
public Acc add(Item value, Acc accumulator) {
accumulator.inc(value.getSummary());
accumulator.addUid(value.getUid);
return accumulator;
}
@Override
public Tuple2<Integer,List<String>> getResult(Acc accumulator) {
List<String> newL = Lists.newArrayList(accumulator.getUids());
accumulator.setUids(Lists.newArrayList());
return Tuple2.of(accumulator.getSum(), newL);
}
@Override
public Acc merge(Acc a, Acc b) {
.....
}
}
在聚合过程函数中,我将列表刷新到状态,如果需要保存到数据库,我将清除状态并在状态中保存标志以指示它。
但在我看来它是歪曲的。而且我不确定这是否适合我。
这种情况有更好的解决办法吗?
【问题讨论】:
-
您的问题令人困惑。对我来说,您似乎需要 2 个查询。一个在窗口时间内求和,另一个是无窗口。
-
没有。他们都需要同一个窗口。 “uids”列表在摘要中为我提供了所需的调试指示
标签: apache-flink flink-streaming