How to add a SparkListener from pySpark in Python?

Viewed 7027

I want to create a Jupyter/IPython extension to monitor Apache Spark Jobs.

Spark provides a REST API.

However instead of polling the server, I want the event updates to be sent through callbacks.

I am trying to register a SparkListener with the SparkContext.addSparkListener(). This feature is not available in the PySpark SparkContext object in Python. So how can I register a python listener to Scala/Java version of the context from Python. Is it possible to do this through py4j? I want python functions to be called when the events fire in the listener.

2 Answers

I know this is a very old question. However I ran into this very same issue where we had to configure a custom developed listener in a PySpark application. Possible that in the last few years the approach changed.

All we had to do is to specify the dependent jar file that contained the listener jar and also set a --conf spark.extraListeners property.

Example

--conf spark.extraListeners=fully.qualified.path.to.MyCustomListenerClass --conf my.param.name="hello world"

MyCustomListenerClass can have a single argument constructor that accepts a SparkConf object. If you want to pass any parameters to your listener, just set them as configuration key-values and you should be able to access them from the constructor.

Example

public MyCustomListenerClass(SparkConf conf) {
        this.myParamName = conf.get("my.param.name", "default_param_value");
}

Hope this helps someone looking for a simpler strategy. The approach works on both Scala and PySpark because nothing changes in the spark application, the framework takes care of registering your listener by just passing the extraListeners parameter.

Related