Local Combiner for GroupByKey on Flink Streaming jobs

2023-05-19 Thread Talat Uyarer via user
Hi, I have a stream aggregation job which is running on Flink 1.13 I generate DAG by using Beam SQL. My SQL query has a TUMBLE window. Basically My pipeline reads from kafka aggregate, counts/sums some values by streamin aggregation and writes a Sink. BeamSQl uses Groupbykey for the aggregation p

Watermark Alignment on Flink Runner's UnboundedSourceWrapper

2023-05-19 Thread Talat Uyarer via user
Hi All, I have a stream aggregation job which reads from Kafka and writes some Sinks. When I submit my job Flink checkpoint size keeps increasing if I use unaligned checkpoint settings and it does not emit any window results. If I use an aligned checkpoint, size is somewhat under control(still bi