with the current API this should do what you are after:

val input = ...
val result = input
  .reduceWindow( /* your reduce function */ )

With the reduce function you should be able to implement any custom aggregations. You can also use foldWindow() if you want to do a functional fold over the window.

I hope this helps.


On Fri, 21 Aug 2015 at 14:51 Philipp Goetze <philipp.goetze@tu-ilmenau.de> wrote:
Hello community,

how do I define a custom aggregate function in Flink Streaming (Scala)?
Could you please provide an example on how to do that?

Thank you and best regards,