spark structured streaming pyspark, applyInPandas being called immediately and not waiting for window to expire

Viewed 139

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

0 Answers
Related