Pyspark - Applying custom function on structured streaming

Viewed 219

I have 4 columns ['clienttimestamp",'sensor_id','actvivity',"incidents"]. From kafka stream, i consume data,preprocess and aggregate in window.

If i do with groupby with ".count()", The stream works very well writing each window with their count in the console.

This works,

df = df.withWatermark("clientTimestamp", "1 minutes")\
                        .groupby(window(df.clientTimestamp, "1 minutes", "1 minutes"), col('sensor_type')).count()
query = df.writeStream.outputMode("append").format('console').start() 
query.awaitTermination()

But the real motive is to find the total time for which critical activity was live. i.e. For each sensor_type, i group the data by window and i get the list of critical activity and find the total time for which the all critical activity lasted" (The code is below). But am not sure if i am using the udf in right way! Because below method does not work! Can anyone provide an example of applying a custom function for each group of window and to write the output to console.

This does not work

@f.pandas_udf(schemahh, f.PandasUDFType.GROUPED_MAP)
def calculate_time(pdf):
    pdf = pdf.reset_index(drop=True)
    total_time = 0
    index_list = pdf.index[pdf['activity'] == 'critical'].to_list()
    for ind in index_list:
        start = pdf.loc[ind]['clientTimestamp']
        end = pdf.loc[ind + 1]['clientTimestamp']
        diff = start - end
        time_n_mins = round(diff.seconds / 60, 2)
        total_time = total_time + time_n_mins
    largest_session_time = total_time
    new_pdf = pd.DataFrame(columns=['sensor_type', 'largest_session_time'])
    new_pdf.loc[0] = [pdf.loc[0]['sensor_type'], largest_session_time]
    return new_pdf

df = df.withWatermark("clientTimestamp", "1 minutes")\
                        .groupby(window(df.clientTimestamp, "1 minutes", "1 minutes"), col('sensor_type'), col('activity')).apply(calculate_time)
query = df.writeStream.outputMode("append").format('console').start() 
query.awaitTermination()
0 Answers
Related