kafka-streams: joining historical data with a stream

Viewed 482

Given that I have the three following topics:

  • topic-a: events are produced continuously - topic is partitioned but not compacted - updates for previous messages are common
  • topic-b: events occur on a monthly basis - topic is partitioned but not compacted - events are related to events in topic a - all keys from topic a will be referenced in this topic
  • topic-c: output topic for events generated by merging a and b

Objective: Messages from topic-b need to be enriched with the events stored in topic-a and then published on topic-c.

Problem: I cannot join the topics directly, as topic-a contains many duplicates and I need to have exactly a 1:1 match between topic-a and topic-b.

Idea 1

  • Create a KTable aTable from topic-a - aggregate events
KStream<String, DomainObjectA> aStream = kStreamBuilder.stream(
                aTopic,
                Consumed.with(Serdes.String(), aSerde));
KTable<String, DomainObjectA> aTable = aStream
               .groupByKey()
               .reduce(
                       (value1, value2) -> value2
               );
  • Join events from topic-b with KTable based on stream a (aTable)
KStream<String, DomainObjectB> bStream = kStreamBuilder.stream(topicB);
KStream<String, DomainObjectC> cStream = bStream.leftJoin(aTable, getJoiner())

Problem: this might work, yet AFAIK, the state-store and the underlying topic created by KTable aTable will grow indefinitely.

Idea 2:

  • Use windowing:
KStream<String, DomainObjectA> aStream = kStreamBuilder.stream(
                aTopic,
                Consumed.with(Serdes.String(), aSerde));
KTable<Windowed<String>, DomainObjectA> aTable = aStream
                .groupByKey()
                .windowedBy(TimeWindows.of(Duration.ofDays(30L)))
                .reduce(
                        (value1, value2) -> value2,
                        Materialized.<String, DomainObjectA, WindowStore<Bytes, byte[]>>as("a-storage").withRetention(Duration.ofDays(30L)));

Problem: If I understand correctly, this approach should ensure that the underlying topic for the KTable and the state-store are bounded and entries will be cleaned up. However, I end up with a KTable<Windowed<String>, DomainObjectA> which I cannot join with my stream:

KStream<String, DomainObjectB> bStream = kStreamBuilder.stream(topicB);
KStream<String, DomainObjectC> cStream = bStream.leftJoin(aTable, getJoiner()); // not possible

I was able to join by turning stream-b into a windowed KTable as well - but I am not sure I like this approach as it will create yet another (this time unnecessary) topic.

Idea 3:

Doing things manually by:

  • Creating a new compacted topic (topic-d) with 30 day retention-period
  • Create a new topology which pushes all messages from topic-a to topic-d
  • Join all values from topic-b with topic-d - directly joining the steams without any KTables While having some overhead, I think this might work, however it feels like I am by-passing the framework somehow.

So in the end my question is: what is the best way to achieve a join of historical data with a stream using kafka-streams? While idea 2 might work I think idea 3 is the better approach as it requires only one additional topic and no state-stores.

0 Answers
Related