Spark Structured Streaming - join 2 dataframes based on condition

Viewed 137

Using spark structured streaming for real time data streaming- I am using a static dataframe to enrich the streaming dataframe. Now that based on a column value I have to join the stream dataframe with either one of two static dataframes.(using 2 join conditions - use either one of these based on the column values)

    val df_join1 = streamDF.as("a").where(streamDF.col("col1") isin ("value1","value2")).join(staticDF.as("b"), joinCondition01, "left_outer")
    
    
    val df_join2 = streamDF.as("a").where(streamDF.col("col1") isin ("value3","value4")).join(staticDF.as("b"), joinCondition02, "left_outer")
    
    val unionDF = df_join1.union(df_join2)
    
    
    val writequery1 = unionDF.selectExpr("to_json(struct(*)) AS value")
                      .writeStream
                      .format("kafka")
                      .outputMode("append")
                      .option("kafka.bootstrap.servers", "serverNames:port")
                      .option("kafka.security.protocol","SSL")
                      .option("kafka.ssl.truststore.location","/location/client_truststore.jks")
                      .option("kafka.ssl.keystore.location", "/location/client_keystore.jks")
                      .option("topic", "topic_test01")
                      .option("checkpointLocation","/checkpointLocation/tmp_01")
                      .start()

    val writequery2 = unionDF.selectExpr("to_json(struct(*)) AS value")
                      .writeStream
                      .format("kafka")
                      .outputMode("append")
                      .option("kafka.bootstrap.servers", "serverNames:port")
                      .option("kafka.security.protocol","SSL")
                      .option("kafka.ssl.truststore.location","/location/client_truststore.jks")
                      .option("kafka.ssl.keystore.location", "/location/client_keystore.jks")
                      .option("topic", "topic_test02")
                      .option("checkpointLocation","/checkpointLocation/tmp_02")
                      .start()    
    
   sparksession.streams.awaitAnyTermination()

But I am getting streamingQueryException - Query [] terminated with exception: assertion failed: There are [1] sources in the checkpoint offsets and now there are [2] sources requested by the query. Cannot continue. Caused by: AssertionError: assertion failed: There are [1] sources in the checkpoint offsets and now there are [2] sources requested by the query. Cannot continue.

Can anyone please help, I am new to spark and structured streaming. I believe there would be a better way out. Please suggest.

0 Answers
Related