Spark Structured Streaming inconsistent output to multiple sinks

Viewed 42

I am using spark structured streaming to read data from Kafka and apply some udf to the dataset. The code as below :

calludf = F.udf(lambda x: function_name(x))

dfraw = spark.readStream.format('kafka') \
    .option('kafka.bootstrap.servers', KAFKA_CONSUMER_IP) \
    .option('subscribe', topic_name) \
    .load()
df = dfraw.withColumn("value", F.col('value').cast('string')).withColumn('value', calludf(F.col('value')))

ds = df.selectExpr("CAST(value AS STRING)") \
    .writeStream \
    .format('console') \
    .option('truncate', False) \
    .start()

dsf = df.selectExpr("CAST (value AS STRING)") \
    .writeStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", KAFKA_CONSUMER_IP) \
    .option("topic", topic_name_two) \
    .option("checkpointLocation", checkpoint_location) \
    .start()

ds.awaitTermination()
dsf.awaitTermination()

Now the problem is that I am getting 10 dataframes as input. 2 of them failed due to some issue with the data which is understandable. The console displays rest of the 8 processed dataframes BUT only 6 of those 8 processed dataframes are written to the Kafka topic using dsf steaming query. Even though I have added checkpoint location to it but it is still not working.

PS: Do let me know if you have any suggestion regarding the code as well. I am new to spark structured streaming so maybe there is something wrong with the way I am doing it.

0 Answers
Related