# Task exception could not be deserialized

**URL:** <https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960>\
**Category:** Elasticsearch\
**Tags:** es-hadoop\
**Created:** [March 31, 2016, 3:16pm UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960 "2016-03-31T15:16:31Z")\
**Posts on this page:** 13\
**Page:** 1

<div class="post-metadata">

**Author:** ![Jonathan\_Spooner](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/jonathan_spooner/32/2672_2.png) [@Jonathan\_Spooner](https://discuss.elastic.co/u/Jonathan_Spooner)\
**Post date:** [March 31, 2016, 3:16pm UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/1 "2016-03-31T15:16:31Z")

</div>

I have a Spark job that runs on my localhost but when run on EMR I'm getting a Warning for

```
WARN ThrowableSerializationWrapper: Task exception could not be deserialized
java.lang.ClassNotFoundException: org.elasticsearch.hadoop.EsHadoopException

```

and error

```
16/03/31 16:32:55 ERROR TaskResultGetter: Could not deserialize TaskEndReason: ClassNotFound with classloader org.apache.spark.util.MutableURLClassLoader@216b55e2

```

• My Spark EMR version is 1.5.2  
• I was using `elasticsearch-spark_2.10-2.2.0-rc1` but I also tried the latest release `elasticsearch-spark_2.10-2.2.0`

My code is simply creating a bunch of dummy documents and inserting them into elasticsearch.

```
def deviceDoc(deviceID: String): Any = {
  // snip because this post is limited to 500 chars
 return map;
}
val deviceRDD = deviceAggDF.map( x=>
         deviceDoc( x(0).toString )
       );
EsSpark.saveToEs(deviceRDD, "parent-child-devices/device", Map("es.mapping.id" -> "_id"))

```

Submitting the job

```
spark-submit \
    --jars ./elasticsearch-spark_2.10-2.2.0-rc1.jar \
    --class com.me.MyJob \
    --packages com.databricks:spark-csv_2.11:1.3.0 \
    ./myproject.jar \

```

The error

```
16/03/31 15:11:12 WARN ThrowableSerializationWrapper: Task exception could not be deserialized
java.lang.ClassNotFoundException: org.elasticsearch.hadoop.EsHadoopException
	at java.net.URLClassLoader$1.run(URLClassLoader.java:366)
	at java.net.URLClassLoader$1.run(URLClassLoader.java:355)
	at java.security.AccessController.doPrivileged(Native Method)
	at java.net.URLClassLoader.findClass(URLClassLoader.java:354)
	at java.lang.ClassLoader.loadClass(ClassLoader.java:425)
	at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:308)
	at java.lang.ClassLoader.loadClass(ClassLoader.java:358)
	at java.lang.Class.forName0(Native Method)
	at java.lang.Class.forName(Class.java:278)
	at org.apache.spark.serializer.JavaDeserializationStream$$anon$1.resolveClass(JavaSerializer.scala:67)
	at java.io.ObjectInputStream.readNonProxyDesc(ObjectInputStream.java:1612)
	at java.io.ObjectInputStream.readClassDesc(ObjectInputStream.java:1517)
	at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:1771)
	at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1350)
	at java.io.ObjectInputStream.readObject(ObjectInputStream.java:370)
	at org.apache.spark.ThrowableSerializationWrapper.readObject(TaskEndReason.scala:167)
	at sun.reflect.GeneratedMethodAccessor44.invoke(Unknown Source)
	at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
	at java.lang.reflect.Method.invoke(Method.java:606)
	at java.io.ObjectStreamClass.invokeReadObject(ObjectStreamClass.java:1058)
	at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:1897)
	at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:1798)
	at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1350)
	at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:1997)
	at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:1921)
	at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:1798)
	at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1350)
	at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:1997)
	at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:1921)
	at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:1798)
	at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1350)
	at java.io.ObjectInputStream.readObject(ObjectInputStream.java:370)
	at org.apache.spark.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:72)
	at org.apache.spark.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:98)
	at org.apache.spark.scheduler.TaskResultGetter$$anon$3$$anonfun$run$2.apply$mcV$sp(TaskResultGetter.scala:108)
	at org.apache.spark.scheduler.TaskResultGetter$$anon$3$$anonfun$run$2.apply(TaskResultGetter.scala:105)
	at org.apache.spark.scheduler.TaskResultGetter$$anon$3$$anonfun$run$2.apply(TaskResultGetter.scala:105)
	at org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:1699)
	at org.apache.spark.scheduler.TaskResultGetter$$anon$3.run(TaskResultGetter.scala:105)
	at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145)
	at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615)
	at java.lang.Thread.run(Thread.java:745)
16/03/31 15:11:12 ERROR TaskResultGetter: Could not deserialize TaskEndReason: ClassNotFound with classloader org.apache.spark.util.MutableURLClassLoader@1cba98ca
```

---

<div class="post-metadata">

**Author:** ![costin](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/costin/32/44950_2.png) [@costin](https://discuss.elastic.co/u/costin)\
**Post date:** [April 5, 2016, 2:59pm UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/2 "2016-04-05T14:59:42Z")

</div>

This error occurs every now and then but I can't reproduce it despite several attempts.

As far as I can tell, the problem lies with Spark. Namely one of the tasks fails (maybe because the data to/from ES is malformed or something) and ES-Hadoop throws the exception higher up the stack. Spark catches this, shuts down the workers and tries to bubble up the exception in the master. However since the job classloader has been closed automatically, the exception cannot be deserialized (despite the class being available in the classpath) which leads to a different error, masking the original one.

In fact, if you look closely through the logs, you should see the initial ESHadoopException stacktrace.

P.S. By the way, the es-hadoop package is available as a spark package as well.  
P.P.S. Using spark 1.6.1 might help

---

<div class="post-metadata">

**Author:** ![Jonathan\_Spooner](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/jonathan_spooner/32/2672_2.png) [@Jonathan\_Spooner](https://discuss.elastic.co/u/Jonathan_Spooner)\
**Post date:** [April 6, 2016, 3:52pm UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/3 "2016-04-06T15:52:13Z")

</div>

I have a feeling this is caused by my mock device document function returning **Any**. If I have this function only return a Map with an \_id attribute the job seems to never have errors.

This might just be me not knowing Scala, I believe the correct return type of this function should specify a Map with the nested types.

```
def deviceDoc(deviceID: String): Any = {
  val map = Map(
    "_id" -> deviceID
    ,
    "device" -> Map(
      "carrier" -> random(Set("AT&T", "Verision")),
      "isp" -> random(Set("isp_214", "isp_215")),
      "model" -> random(Set("4", "5s")),
      "os" -> random(Set("iOS", "android"))
    )
    ,
    "demographics" -> Seq(
      Map(
        "age_range" -> random(Set("1-10", "11-20", "21-30", "31-40", "41-50", "old as dirt")),
        "conf_score" -> 57,
        "ethnicity" -> random(Set("asian", "white", "other")),
        "gender" -> random(Set("M", "F")),
        "hhi" -> random(Set("50K-60K", "70K-140K", "150K-160K")),
        "source" -> "datasrc_5"
      )
    )
  );
  
  return map
}
```

---

<div class="post-metadata">

**Author:** ![costin](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/costin/32/44950_2.png) [@costin](https://discuss.elastic.co/u/costin)\
**Post date:** [April 6, 2016, 3:54pm UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/4 "2016-04-06T15:54:07Z")

</div>

Fair enough but the exception should still be able to propagate properly - can you run a basic integration test in a local spark cluster? Do you see the exception being raised and if so, what is it?

---

<div class="post-metadata">

**Author:** ![Jonathan\_Spooner](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/jonathan_spooner/32/2672_2.png) [@Jonathan\_Spooner](https://discuss.elastic.co/u/Jonathan_Spooner)\
**Post date:** [April 6, 2016, 3:58pm UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/5 "2016-04-06T15:58:26Z")

</div>

Actually I take that back. I just ran a job with just the \_id attribute and the error was raised ☹

---

<div class="post-metadata">

**Author:** ![Jonathan\_Spooner](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/jonathan_spooner/32/2672_2.png) [@Jonathan\_Spooner](https://discuss.elastic.co/u/Jonathan_Spooner)\
**Post date:** [April 6, 2016, 4:07pm UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/6 "2016-04-06T16:07:15Z")

</div>

Geez even removing the function and simply returning a Map throws the error.

```
val deviceRDD = aggDF.map(x=> Map(
         "_id" -> x(0)
));
EsSpark.saveToEs(deviceRDD, "parent-child-devices/device", Map("es.mapping.id" -> "_id"))

```

error

```
16/04/06 16:04:47 INFO TaskSetManager: Starting task 29.1 in stage 1.0 (TID 340, ip-10-136-126-99.ec2.internal, PROCESS_LOCAL, 2316 bytes)
16/04/06 16:04:47 WARN ThrowableSerializationWrapper: Task exception could not be deserialized
java.lang.ClassNotFoundException: org.elasticsearch.hadoop.EsHadoopException
	at java.net.URLClassLoader$1.run(URLClassLoader.java:366)
	at java.net.URLClassLoader$1.run(URLClassLoader.java:355)
	at java.security.AccessController.doPrivileged(Native Method)
	at java.net.URLClassLoader.findClass(URLClassLoader.java:354)
	at java.lang.ClassLoader.loadClass(ClassLoader.java:425)
	at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:308)
	at java.lang.ClassLoader.loadClass(ClassLoader.java:358)
	at java.lang.Class.forName0(Native Method)
	at java.lang.Class.forName(Class.java:278)
	at org.apache.spark.serializer.JavaDeserializationStream$$anon$1.resolveClass(JavaSerializer.scala:67)
	at java.io.ObjectInputStream.readNonProxyDesc(ObjectInputStream.java:1612)
	at java.io.ObjectInputStream.readClassDesc(ObjectInputStream.java:1517)
	at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:1771)
	at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1350)
	at java.io.ObjectInputStream.readObject(ObjectInputStream.java:370)
	at org.apache.spark.ThrowableSerializationWrapper.readObject(TaskEndReason.scala:167)
	at sun.reflect.GeneratedMethodAccessor29.invoke(Unknown Source)
	at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
	at java.lang.reflect.Method.invoke(Method.java:606)
	at java.io.ObjectStreamClass.invokeReadObject(ObjectStreamClass.java:1058)
	at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:1897)
	at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:1798)
	at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1350)
	at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:1997)
	at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:1921)
	at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:1798)
	at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1350)
	at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:1997)
	at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:1921)
	at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:1798)
	at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1350)
	at java.io.ObjectInputStream.readObject(ObjectInputStream.java:370)
	at org.apache.spark.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:72)
	at org.apache.spark.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:98)
	at org.apache.spark.scheduler.TaskResultGetter$$anon$3$$anonfun$run$2.apply$mcV$sp(TaskResultGetter.scala:108)
	at org.apache.spark.scheduler.TaskResultGetter$$anon$3$$anonfun$run$2.apply(TaskResultGetter.scala:105)
	at org.apache.spark.scheduler.TaskResultGetter$$anon$3$$anonfun$run$2.apply(TaskResultGetter.scala:105)
	at org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:1699)
	at org.apache.spark.scheduler.TaskResultGetter$$anon$3.run(TaskResultGetter.scala:105)
	at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145)
	at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615)
	at java.lang.Thread.run(Thread.java:745)
16/04/06 16:04:47 ERROR TaskResultGetter: Could not deserialize TaskEndReason: ClassNotFound with classloader org.apache.spark.util.MutableURLClassLoader@dd37992
16/04/06 16:04:47 WARN TaskSetManager: Lost task 155.1 in stage 1.0 (TID 265, ip-10-169-62-32.ec2.internal): UnknownReason
```

---

<div class="post-metadata">

**Author:** ![Jonathan\_Spooner](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/jonathan_spooner/32/2672_2.png) [@Jonathan\_Spooner](https://discuss.elastic.co/u/Jonathan_Spooner)\
**Post date:** [April 6, 2016, 4:11pm UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/7 "2016-04-06T16:11:28Z")

</div>

Not sure how to run an integration test on the cluster. I am able to run spark-shell and walk through each step.

---

<div class="post-metadata">

**Author:** ![Jonathan\_Spooner](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/jonathan_spooner/32/2672_2.png) [@Jonathan\_Spooner](https://discuss.elastic.co/u/Jonathan_Spooner)\
**Post date:** [April 6, 2016, 4:21pm UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/8 "2016-04-06T16:21:56Z")

</div>

The difficult thing is that the error doesn't consistently happen. I'm running the job multiple times on a small file and I seem to get the error 50% of the time.

---

<div class="post-metadata">

**Author:** ![Jonathan\_Spooner](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/jonathan_spooner/32/2672_2.png) [@Jonathan\_Spooner](https://discuss.elastic.co/u/Jonathan_Spooner)\
**Post date:** [April 6, 2016, 6:07pm UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/9 "2016-04-06T18:07:44Z")

</div>

So I'm stepping through my job in spark-shell to ensure the data in my RRD is ok. I can not identify any issues with the data.

---

<div class="post-metadata">

**Author:** ![Jonathan\_Spooner](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/jonathan_spooner/32/2672_2.png) [@Jonathan\_Spooner](https://discuss.elastic.co/u/Jonathan_Spooner)\
**Post date:** [April 6, 2016, 9:30pm UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/10 "2016-04-06T21:30:17Z")

</div>

Looks like my EMR cluster was using an old version of Spark. I have the a cluster running with 1.6 and I believe it is fixed. Thanks!

---

<div class="post-metadata">

**Author:** ![costin](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/costin/32/44950_2.png) [@costin](https://discuss.elastic.co/u/costin)\
**Post date:** [April 7, 2016, 9:29am UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/11 "2016-04-07T09:29:51Z")

</div>

So upgrading the Spark cluster to 1.6 seems to have fixed the problem, right? If the issue keeps popping up, please speak up.

---

<div class="post-metadata">

**Author:** ![larghir](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/larghir/32/61526_2.png) [@larghir](https://discuss.elastic.co/u/larghir)\
**Post date:** [March 23, 2017, 10:27am UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/12 "2017-03-23T10:27:13Z")

</div>

I'm having the same issue with a long-running job, a few tasks are failing and am seeing the same error stack. Running Spark 1.5.2.  
Is updating to Spark 1.6.x guaranteed to fix this?  
Thanks

---

<div class="post-metadata">

**Author:** ![system](https://us1.discourse-cdn.com/elastic/original/3X/1/a/1ac57faf039f6b580b3f104ef42a2a89e41014de.png) [@system](https://discuss.elastic.co/u/system)\
**Post date:** [July 6, 2017, 1:22pm UTC](https://discuss.elastic.co/t/task-exception-could-not-be-deserialized/45960/13 "2017-07-06T13:22:24Z")

</div>


