Event data aggregation using Kafka Streams in a system with high volume of data

Viewed 46

In our company we have a use which can be mapped to the one exposed in Confluent about joining movies and ratings using a KStream and KTable.

Here is the problem statement described in Confluent article:

Suppose you have a set of movies that have been released and a stream of ratings from moviegoers about how entertaining they are. In this tutorial, we'll write a program that joins each rating with content about the movie. Related pattern: Event Joiner

The short answer they provide is the following:

KStream<String, Rating> ratings = ...
KTable<String, Movie> movies = ...
final MovieRatingJoiner joiner = new MovieRatingJoiner();
KStream<String, RatedMovie> ratedMovie = ratings.join(movies, joiner);

Source: https://developer.confluent.io/tutorials/join-a-stream-to-a-table/kstreams.html

However, in our case the data we are dealing with has the following properties:

  1. The number of total movies is about 500 million, while the number of ratings is even higher.
  2. We always need to have the complete movie catalog in our lookup store, otherwise we may receive a rating for a movie we don't have and the join will not work.
  3. A new rating can be received anytime, and a new movie can be received anytime.
  4. We just need to send the rated movie event (join result) if the rating was done after one month ago. If the rating was done more than one month ago we don't want to send the rated movie event anymore.

Then, considering properties P1 and P2, I think for movies lookup is more appropriate to use a persistent database instead of a KTable. KTable can use local disk storage and is fault-tolerant, since it can be rebuilt using the kafka topic in case it crashes. Even though, that process could take a long time if we are talking about millions of movies.

Having P3, we need to guaranty there is no race condition which may lead to an event being lost. With P4 we are considering a time window of one month.

As you can imagine, most likely, the rating will belong to an already existent movie, so the initial lookup will be successful more than 95% of the times. However, it is possible to receive a rating for a new movie which is not yet present in our system, or a movie which we will never receive.

Apart from KStream-KTable, I was also thinking on a KStream-KStream windowed join, but considering the volume of the data and the window of one month, I expect the memory usage to be too high.

Based on the properties described above, is there any recommended pattern using KStream-KStream or KStream-KTable joins which could fit this use case?

Otherwise I believe the best option we have is to use two persistent datasources for the lookup, like this: enter image description here

0 Answers
Related