spark structured streaming is not allowing Window function to perform lag,lead operations. So I am trying to use applyInpandas function.
I have a tumbling window of 5 minutes with watermark set to 1 minute in append mode.I need to wait till window expires and apply my UDF function on it.So I am using applyInpandas.
The problem is my custom UDF function is being called immediately after data arrives and not waiting till the window expires.
My code
schema = StructType([
StructField("date", TimestampType(), True),
StructField("id", StringType(), True),
StructField("x1", IntegerType(), True),
StructField("y1", IntegerType(), True),
])
newschema = StructType([
StructField("date", TimestampType(), True),
StructField("id", StringType(), True),
StructField("x1", DoubleType(), True),
StructField("y1", DoubleType(), True),
StructField("x2", DoubleType(), True),
StructField("y2", DoubleType(), True),
StructField("dis", DoubleType(), True)
])
spark = SparkSession.builder.appName("testing").getOrCreate()
spark.sparkContext.setLogLevel('WARN')
truck_data = spark.readStream.format("kafka")\
.option("kafka.bootstrap.servers", KAFKA_BROKER)\
.option("subscribe", KAFKA_INPUT_TOPIC)\
.option("startingOffsets", "latest") \
.load()
truck_data = truck_data.withColumn("value", col("value").cast(StringType()))\
.withColumn("jsonData", from_json(col("value"), schema)) \
.select("JsonData.*")\
def f(pdf):
print("inside applyinpandas")
pdf = pdf.sort_values(by='date')
pdf = pdf.assign(x2=pdf.x1.shift(1))
pdf = pdf.assign(y2=pdf.y1.shift(1))
pdf = pdf.assign(dis=np.sqrt((pdf.x1 - pdf.x2)**2 + (pdf.x1 - pdf.x2)**2))
pdf = pdf.fillna(0)
return pdf
truck_data = truck_data.withWatermark("date", "1 minute").\
.groupBy(window("date", "5 minute"), "id")\
.applyInPandas(f, schema=newschema)
query = truck_data\
# .writeStream\
# .format("console") \
# .option("numRows", 300)\
# .option("truncate", False)\
# .start().awaitTermination()
Sample input data
{"date":"2022-03-23 09:04:32.242637","id":"B","x1":3,"y1":3}
{"date":"2022-03-23 09:04:32.242737","id":"A","x1":2,"y1":2}
{"date":"2022-03-23 09:04:29.242737","id":"A","x1":1,"y1":1}
{"date":"2022-03-23 09:04:55.242737","id":"A","x1":6,"y1":6}
{"date":"2022-03-23 09:04:40.242737","id":"B","x1":7,"y1":7}
{"date":"2022-03-23 09:04:29.242737","id":"B","x1":1,"y1":1}
{"date":"2022-03-23 09:04:44.242737","id":"A","x1":5,"y1":5}
{"date":"2022-03-23 09:04:35.242737","id":"B","x1":5,"y1":5}
{"date":"2022-03-23 09:04:35.242737","id":"A","x1":3,"y1":3}
{"date":"2022-03-23 09:04:40.242737","id":"A","x1":4,"y1":4}
{"date":"2022-03-23 09:04:44.242737","id":"B","x1":9,"y1":9}
{"date":"2022-03-23 09:04:55.242737","id":"B","x1":11,"y1":11}
{"date":"2022-03-23 09:06:55.242737","id":"B","x1":11,"y1":11}
Output required output image from pycharm