How to use Spark Accumulators to instrument Pyspark Jobs using SparkListener?

Viewed 283

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

  1. Is there something I am missing when using pyspark accumulators and is there a way to get them to show up in sparklistener?
  2. If not can I some how pickle the scala spark object which essentially is a py4j JavaObject in the method above?
0 Answers
Related