Pyspark structured streaming trigger=availableNow get stuck on occasion

Viewed 22

I have several tasks of streaming tables in pyspark (running on Databricks).

For the most part, it looks something like this:

stream = (
    spark
    .readStream
    .option("maxFilesPerTrigger", 20)
    .table("my_table")
)
transformed_df = stream.transform(some_function)

 _ = (
    transformed_df
    .writeStream
    .trigger(availableNow=True)
    .outputMode("append")
    .option("checkpointLocation", "/path/to/_checkpoints/output_table/input_table/")
    .foreachBatch(lambda df, epochId: batch_writer(df, epochId, "output_table"))
    .start()
    .awaitTermination()
)

This works as expected and I can see each batch commit and offset in the checkpoints path. The problem is that, on occasion, the stream writer will just carry on (no new data coming in) and I can see the commit and offset go into the thousands. If I stop the job, it does not do it again and then later again.

I am running Pyspark on Databricks (Apapche Spark 3.2.1, Scala 2.12).

Should I do a check like if len(transformed_df.take(1)) > 0: # write stream? I suspect this will almost always be true and the stream is only determined at write.

0 Answers
Related