Jad Naous
08/23/2023, 8:58 PMDatastreamUtils.reinterpretAsKeyedStream() and using the current task index (mapped in a rich function) as the key for the local aggregation. I'm not really sure that's the right thing to do, or what the behavior will end up looking like in case of restarts since the state for one partition may end up at some other task index later. What's the best way to do a local aggregation in Flink?Jad Naous
08/23/2023, 11:38 PM