Pyspark streaming low cpu utilization

Viewed 38

I'm trying to run some processing on a stream with Spark streaming in python. The stream is giving out 2000 datapoints each second, and the processing is supposed to give results each 4 seconds (slideDuration=4), but in fact it is being delayed a lot (printing the each result after ~30 seconds).

from pyspark import SparkContext
from pyspark.streaming import StreamingContext

# Create a local StreamingContext with two working thread and batch interval of 1 second
sc = SparkContext("local[2]", "NetworkWordCount")
ssc = StreamingContext(sc, 2)
ssc.checkpoint("./checkpoints")  # set checkpoint directory

lines = ssc.socketTextStream("localhost", 9000)
lines = lines.map(lambda x : (x, 1/40000))
lines = lines.reduceByKeyAndWindow(lambda a, b: a + b, lambda a, b: a - b, windowDuration=20, slideDuration=4)
lines = lines.filter(lambda x : x[1] > 0.03)

lines.pprint()
ssc.start()
ssc.awaitTermination()

The cpu utilisation though is around 40% and does not get higher.

0 Answers
Related