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
KTableaTable 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
KTablebased 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.