Spark Streaming Job starts reading from beginning instead of where it stopped consuming

Viewed 90

I am trying to use Dataproc on Google Cloud Platform for my Spark Streaming jobs.

I use Kafka as my source and try to write it to MongoDB. Its working fine, but after the job fails it starts to read the messages from my Kafka topic from the beginning instead of from where it stopped.

Here is my config for reading from Kafka:

  clickstreamTestDf = (
  spark
  .readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", confluentBootstrapServers)
  .option("kafka.security.protocol", "SASL_SSL")
  .option("kafka.sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username='{}' password='{}';".format(confluentApiKey, confluentSecret))
  .option("kafka.ssl.endpoint.identification.algorithm", "https")
  .option("kafka.sasl.mechanism", "PLAIN")
  .option("subscribe", "customer_experience")
  .option("failOnDataLoss", "false")
  .option("startingOffsets", "earliest")
  .load()
)

And here is my write stream code:

finished_df.writeStream \
    .format("mongodb")\
    .option("spark.mongodb.connection.uri", connectionString) \
    .option("spark.mongodb.database", "Company-Environment") \
    .option("spark.mongodb.collection", "customer_experience") \
    .option("checkpointLocation", "gs://firstsparktest_1/checkpointCustExp") \
    .option("forceDeleteTempCheckpointLocation", "true") \
    .outputMode("append") \
    .start() \
    .awaitTermination()
  1. Do I need to set startingOffsets to latest? I tried but it still didn't read from where it stopped.
  2. Can I use checkpointLocation like this? Is it okay to use a directory in google storage?
  3. I want to run the streaming job, stop it, delete the Dataproc cluster and then create a new one the next day and continue reading from where it left off. Is that possible and how?

Really need some help here!

0 Answers
Related