Structured Streaming runs out of memory after running for a while even with spark.catalog.clearCache()

Viewed 252

spark version: 2.4.0

My application is a simple pyspark-kafka structured streaming application using park-sql-kafka-0-10_2.11:2.4.0

To avoid this possible memory leak problem , I am also using foreachBatch for each microbatch.

Also, each microbatch is supposed to be composed of <= 1000 rows meaning, it is very unlikely to cause the out of memory issue as long as caches are cleared properly. To be extra cautious, I called spark.catalog.clearCache() at the end of each microbatch to ensure all caches are cleared.

However, after having it run for a while (~ 30 mins) it raises the following issue.

22/01/11 10:39:36 ERROR Client: Application diagnostics message: Application application_1641893608223_0002 failed 2 times due to AM Container for appattempt_1641893608223_0002_000002 exited with  exitCode: -104
Failing this attempt.Diagnostics: Container [pid=17995,containerID=container_1641893608223_0002_02_000001] is running beyond physical memory limits. Current usage: 1.4 GB of 1.4 GB physical memory used; 4.4 GB of 6.9 GB virtual memory used. Killing container.

Even though 1.4 GB is a small amount of memory, each microbatch itself is pretty small as well so it shouldn't be a problem.

Also, there are a lot of tasks stacked in the Kafka-Q, In order to prevent the overload in the spark streaming, I have set spark.streaming.blockInterval to 40000ms and maxOffsetsPerTrigger to 10.

What could be possibly causing this out-of-memory issue?

0 Answers
Related