I am consuming data as structured stream from multiple topics with Pyspark. I have a value column and in that column I have json like string: {"message":"abc", "metrics":{"metric1":"abc", "metric2":123, "metric3":"01/01/2022 00:00:00"}} . I am getting some specific keys (columns) from metrics column this way:
value_schema = StructType([StructField("metrics", StringType(), True)])
topic1_schema = StructType([StructField("metrics1", StringType(), True),
...
StructField("metricsN", StringType(), True)])
topic1_raw = spark \
.readStream \
...
.load() \
.selectExpr("CAST(value AS STRING)") \
.withColumn('value', from_json(col('value'), value_schema)) \
.withColumn('metrics', from_json(col('value.metrics'), topic1_schema)) \
.select(col('metrics.*'))
topic1_batch = topic1_raw \
.writeStream \
.outputMode("append") \
.format("console") \
.start()
topic2_schema = StructType([StructField("metricsA", StringType(), True),
...
StructField("metricsZ", StringType(), True)])
topic2_raw = spark \
.readStream \
...
.load() \
.selectExpr("CAST(value AS STRING)") \
.withColumn('value', from_json(col('value'), value_schema)) \
.withColumn('metrics', from_json(col('value.metrics'), topic2_schema)) \
.select(col('metrics.*'))
topic2_batch = topic2_raw \
.writeStream \
.outputMode("append") \
.format("console") \
.start()
topic1_batch.awaitTermination()
topic2_batch.awaitTermination()
When I run above code separately for each topic I can get my expecting result something like that:
+--------------------+-------+--------+---------+
| metric1|metric3| metric7| metricN|
+--------------------+-------+--------+---------+
|01/01/2022 00:00:...| abc...|12345678| 0|
+--------------------+-------+--------+---------+
+-------+--------+---------+
|metricA| metricB| metricC|
+-------+--------+---------+
| abc...|01/01/22| 123...|
+-------+--------+---------+
But when I run the whole code together, meaning streaming/writing data from all topics I have some values missed:
+--------------------+-------+--------+---------+
| metric1|metric3| metric7| metricN|
+--------------------+-------+--------+---------+
|01/01/2022 00:00:...| abc...|12345678| ***|
+--------------------+-------+--------+---------+
+-------+--------+---------+
|metricA| metricB| metricC|
+-------+--------+---------+
| abc...| null | 123...|
+-------+--------+---------+
What might be the possible reason for that and how to fix it?