I am trying to use accumulators in pyspark to instrument my udf's or custom spark methods in my pyspark jobs. I have written a custom SparkListener in java/scala to listen to the accumulator values onStageComplete. Here the SparkListener I have
public class MyListener extends SparkListener {
@Override
public void onStageCompleted(SparkListenerStageCompleted stageCompleted) {
final Map<Object, AccumulableInfo> accumulableInfoMap = JavaConverters.
mapAsJavaMapConverter(stageCompleted.stageInfo().accumulables())
.asJava();
accumulableInfoMap.forEach((k, accInfo) -> {
final String name = accInfo.name().getOrElse(null);
final String nonNullName = Optional.ofNullable(name).orElse("");
final Optional<Double> accumResult = Optional.ofNullable(accInfo.value()
.getOrElse(null))
.map(value -> {
Double val;
try {
val = Double.valueOf(value.toString());
} catch (Exception e) {
val = null;
}
return val;
});
System.out.println("Printing Accumulator");
System.out.println("name:" + name + " value:" + accumResult.orElse((0.0d))); }
);
}
}
This works well when you are writing scala code and instrumenting using longAccumulator. However seems like named accumulators have not made their way into pyspark yet.
When I use pyspark accumulators in a pyspark shell which i spin up using the following command pyspark --driver-class-path MyJar-1.0.jar --conf spark.extraListeners=package.subpackage.MyListener the pyspark accumulators are not picked up by the sparklistener. MyJar-1.0.jar contains MyListener class I have implemented.
I am using the following test code in my pyspark shell.
def filter_non_42(item, accumulator):
if item % 2 == 0:
accumulator += 1
return '42' in str(item)
from functools import partial
accumulator = sc.accumulator(0)
counting_filter = partial(filter_non_42, accumulator=accumulator)
sc.range(0, 10000).filter(counting_filter).sum()
the output I get from my sparkListener is as follows
Printing Accumulator
name:internal.metrics.executorDeserializeTime value:794.0
Printing Accumulator
name:internal.metrics.executorCpuTime value:1.63990621E8
Printing Accumulator
name:internal.metrics.executorRunTime value:1066.0
Printing Accumulator
name:internal.metrics.jvmGCTime value:182.0
Printing Accumulator
name:internal.metrics.diskBytesSpilled value:0.0
Printing Accumulator
name:internal.metrics.memoryBytesSpilled value:0.0
Printing Accumulator
name:internal.metrics.executorDeserializeCpuTime value:3.01491783E8
Printing Accumulator
name:internal.metrics.resultSize value:2928.0
49995000
As seen above it prints all of native spark accumulators but none of my custom accumulators. I have tried going through the pyspark SparkContext code(context.py) and from what I understand Accumulators in pyspark are independent to those in scala spark.
Trying to work around this I tried to reaching for the java spark context object from pyspark and get a hold of the java long accumulator as follows.
acc = sc._jsc.sc().longAccumulator("MyAccomulator")
The above works well, however pyspark has a problem trying to pickle this accumulator when used similarly in filter_non_42 function above.
My question is
- Is there something I am missing when using pyspark accumulators and is there a way to get them to show up in sparklistener?
- If not can I some how pickle the scala spark object which essentially is a py4j JavaObject in the method above?