Null in output when reading/writing structured streams from multiple topics in Pyspark

Viewed 21

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?

0 Answers
Related