Problem Formulation
I am currently developing an application using Spark Structured Streaming through the PySpark API. I am hosting this application in an AWS EMR cluster.
The main flow of my application is the following:
- Read from a Kafka Topic
- Transform each micro-batch, using forEachBatch, as follows: 2. Perform simple and stateless aggregations (groupBy + count). 4. Write the micro-batch to a Kafka topic
While the app is running, I observed a constant increase in host's CPU and memory usage, under constant load, as seen in the following graphs.

This behavior is consistent across runs. When restarting the application, CPU usage falls to the starting point and then consistently increases.
Questions I have read about the possibility that Garbage Collector is responsible for this increase, but I believe that there is a lot of memory space available.
I plan to add complexity in my application logic, so I am worried that this behavior will have a severe impact on the application's runtime.
My questions are: What is the cause behind this behavior? Can I configure Spark to better handle my usecase?