Rishi045
New Contributor III

Using below code :

conf = {}
df = spark.readStream.format("eventhubs").options(**conf).load()
dataDF = df.select(col("body").cast("STRING"))
data = dataDF.select(json_tuple(col("body"),"table","op_type","records","op_ts")) \
            .toDF("table","op_type","records","op_ts")
 
final_data = data.withColumn("records_json",from_json(col("records"),reqSchema))
final_data = final_data.select(
    *[col("records_json." + field).alias(field) for field in reqSchema.fieldNames()],
    col("op_type"),
    col("op_ts"))
final_data.orderBy(col("op_ts").desc())
final_data = final_data.dropDuplicates([primaryKey])  
final_data = final_data.distinct()
final_data = final_data.drop(final_data.op_ts)
final_data = final_data.drop(final_data.op_type)
 
final_data.coalesce(1).writeStream \
    .format("parquet") \
    .outputMode("append") \
    .option("checkpointLocation",checkPoint_url) \
        .trigger(once=True)\
    .start(rawFilePath_url) \
    .awaitTermination()