I want to set config where my application keep track of consumed messages from kafka. So that whenever it gets failed, it could pick from the last commit or consumed offset onwards.
readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,host2:port2")
.option("subscribe", "topic1")
.load()
.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,host2:port2")
.option("topic", "topic1")
.trigger(Trigger.Continuous("1 second")) // only change in query
.start();
I've read online that checkpointlocation property can be set which can used by spark to keep track of offsets.
Want to know where I can set this property ? can I set in above code within option ? May I know how can I set it properly.
secondly, I'm not able to understand trigger(Trigger.Continuous("1 second")) property. Docs says continuous processing engine will record the progress of the query every second, what kind of progress it record while reading messages from kafka ?