I created a purchases table using Kafka connector (using Flink SQL):
CREATE TABLE purchases (
country STRING,
product STRING
) WITH (
'connector' = 'kafka',
'topic' = 'purchases',
'properties.bootstrap.servers' = 'kafka:29092',
'value.format' = 'json',
'properties.group.id' = '1',
'scan.startup.mode' = 'earliest-offset'
);
and then I perform the following aggregation in Apache Flink:
SELECT `country`, `product`, `purchases`
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY country ORDER BY `purchases` DESC) AS row_num
FROM (select country, product, count(*) as `purchases` from purchases group by country, product))
WHERE row_num <= 3;
There're 2 calculations essentially above:
- grouping of purchases by country and product (
select country, product, count(*) aspurchasesfrom purchases group by country, product) - top-n aggregation of getting top 3 products per country.
I wonder what happens on each new row in purchases table: does Flink recalculate each the grouping for all countries and products and recalculates top 3 product per each country or Flink is smart enough to only recalculate the group to which the new record belongs to and also is smart enough to recalculate top 3 row only for the group to which the product belongs to?
I assume Flink is smart enough in both cases because the overhead to recalculating all of the rows would be very high, however I couldn't find any explicit information regarding this in the documentation.