Options
- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
08-03-2023 04:43 AM
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()