【问题标题】:Flink - how to aggregate in stateFlink - 如何在状态中聚合
【发布时间】: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


【解决方案1】:

在丰富的函数中使用状态。在您的状态和窗口触发刷新值时继续添加uid。官方文档的这个页面有一个例子。

https://ci.apache.org/projects/flink/flink-docs-release-1.11/dev/stream/state/state.html#using-keyed-state

对于您的情况,ListState 会很好用。

编辑:

上述解决方案适用于非窗口情况。对于窗口情况,只需使用具有丰富窗口功能的应用功能的聚合

【讨论】:

  • 谢谢,实际上我需要它在窗口中,所以上面的例子不是我的解决方案,但是你给我带来了方向,只是使用了低级聚合函数 - 应用
  • 很高兴它能帮到你!
  • 好吧,我觉得这对我不利,因为 window 中的 apply 方法的工作方式类似于将所有元素保存在 window 中而不像聚合器那样丢弃它们的过程,所以它会占用我的内存跨度>
  • 您也可以在窗口触发时丢弃您的状态元素。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-05-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多