使用RichAggregateFunction时出现Flink错误

vic*_*tim 4 java apache-flink flink-streaming

我正在尝试在Flink中使用抽象RichAggregateFunction的实现。我希望它是“丰富的”,因为我需要将某些状态存储为聚合器的一部分,并且可以执行此操作,因为我可以访问运行时上下文。我的代码如下所示:

stream.keyBy(...)
   .window(GlobalWindows.create())
   .trigger(...)
   .aggregate(new MyRichAggregateFunction());
Run Code Online (Sandbox Code Playgroud)

但是,我得到一个UnsupportedOperationException说

此聚合函数不能是RichFunction。

我显然没有正确使用RichAggregateFunction。有任何如何正确使用它的示例吗?还是应该将ProcessFunction用于此类操作?

谢谢

Fab*_*ske 6

这不是您的错误。

Flink不支持RichAggregateFunction在组窗口中扩展的功能。

  • 是否还有其他选项可以像RichFunctions中那样在窗口上实现带有状态变量的聚合? (4认同)