We are trying to implement use case which described below, we are having implementation issues which we looking to overcome,
Use Case,
We are trying to do KStream join between 2 Kafka Topics by matching KEY present in messages(JSON) of both streams. Also we should maintain the message sequence as it is arrived in KStream from source.
Scenario is, If Matching Key is yet to arrive in any one of the stream, We should Stop or Retry join until expected key arrives in other topic. We thought to put unmatched records back to KStream but in this case sequence not guaranteed.
Issue 1: How to stop or hold join until the expected key to be arrived in other topic. Eg, KTable has Key 100, But KStream yet to receive Key 100 then We should retry Join or hold KStream until Key 100 arrives.
Issue 2: Is there any way to put Delay or Interval in KStream (Delayed KStream) to receive messages with delayed time or interval.
Additionally we have to build Keyed KStream from Non Keyed Topic (Key will be set by extracting it from Message - JSON)
Java is preferable as We done POC to Join between KTable and KStream
KTable<String, String> leftStream = builder.table("stream1");
KStream<String, String> rightStream = builder.stream("stream2");
KStream<String, String> outstream = rightStream.leftJoin(leftStream, (orig_msg, description) -> {
String new_msg = "";
if (description != null) {
new_msg = orig_msg+"-->Matched--"+description;
}else {
new_msg = orig_msg+"-->UnMatched<--"+description;
}
return new_msg;
});