Hi, you can probably solve this with timers and state, but I am not sure that this the simplest way.
Some guidance how timers and state work: https://beam.apache.org/blog/2017/02/13/stateful-processing.html https://beam.apache.org/blog/2017/08/28/timely-processing.html Regards, Csabi On Thu, 14 Feb 2019 at 20:09, Jeff Klukas <[email protected]> wrote: > I'm not aware that there's currently any way to trigger based on data > size. As you state, AfterPane.elementCountAtLeast lets you trigger based on > number of elements, but from my reading of the implementations of triggers > in the Java SDK, triggers don't have access to sufficient data to maintain > sum of element bytes. > > If you look at the implementation of AfterPaneStatMachine.java [0], the > AfterPane trigger maintains its count via an onElement hook that increments > a piece of state attached to the trigger context. That hook and context > have only limited information available, such as the timestamp of the > element. It does not have access to the element itself and there may not be > a generic method for it to be able to calculate the bytes associated with > an input or output element even if it did have access to the element itself. > > I'm newly reading the code for triggering internals, though, so the > limitations above are somewhat speculation. > > [0] > https://github.com/apache/beam/blob/v2.10.0/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterPaneStateMachine.java#L58-L61 > > On Thu, Feb 14, 2019 at 7:23 AM Rajneesh Agrawal < > [email protected]> wrote: > >> Hi, >> >> I am working on pushing data to s3 from different kafka topics. My >> use-case is to push data in s3 every 15 mins or whenever window size >> exceeds X MBs. What is the simplest way to do this? >> >> I am able to push data every 15 mins using windowing but I am not able to >> trigger the push function when window size exceeds X MBs. Triggering write >> based on number of events is not the correct parameter here since messages >> in kafka are of different size. Does anyone in the group has similar >> use-case? Any help is highly appreciated. >> >> >> >> -- >> Thanks, >> Rajneesh >> >
