# ES-Hadoop PySpark error

**URL:** <https://discuss.elastic.co/t/es-hadoop-pyspark-error/110191>\
**Category:** Elasticsearch\
**Tags:** es-hadoop\
**Created:** [December 4, 2017, 5:42pm UTC](https://discuss.elastic.co/t/es-hadoop-pyspark-error/110191 "2017-12-04T17:42:22Z")\
**Posts on this page:** 3\
**Page:** 1

<div class="post-metadata">

**Author:** ![adesewa](https://avatars.discourse-cdn.com/v4/letter/a/e8c25b/32.png) [@adesewa](https://discuss.elastic.co/u/adesewa)\
**Post date:** [December 4, 2017, 5:42pm UTC](https://discuss.elastic.co/t/es-hadoop-pyspark-error/110191/1 "2017-12-04T17:42:23Z")

</div>

Hello,

I am currently working on a project where I do some fuzzy matching on data in an elasticsearch index. Because I have millions of data in a python dataframe to match against millions of data in the elasticsearch index, it is taking quite a long while as I am having to go through each record in the dataframe.  
So, I thought of using spark's distributed computing power. My aim is this: spark splits the dataframe into different executors and each executor queries the elasicsearch index. This should increase the overall speed. Hence, my work with ES-Hadoop.

I have downloaded the binaries and I ran the following but got the error below. Any idea what I am getting wrong?  
./bin/pyspark --driver-class-path=/Users/xx/ES-Hadoop/elasticsearch-hadoop-6.0.0/dist  
Welcome to  
\_\_\_\_ \_\_  
/ **/** \_\_\_ _**/ /  
 \ / \_ / \_ `/ \_\_/ '/  
/** / ._\_/\_,_/_/ /_/\_\ version 2.2.0  
/_/

Using Python version 2.7.13 (default, Dec 18 2016 07:03:39)  
SparkSession available as 'spark'.

> > > conf = {"es.nodes":"[http://xxxxx.net](http://xxxxx.net)","es.port":9223,"es.resource":"client\_index\_multilang"}  
> > > rdd = sc.newAPIHadoopRDD("org.elasticsearch.hadoop.mr.EsInputFormat", "org.apache.hadoop.io.NullWritable", "org.elasticsearch.hadoop.mr.LinkedMapWritable", conf=conf)  
> > > Traceback (most recent call last):  
> > > File "", line 1, in   
> > > File "/usr/local/Cellar/apache-spark/2.2.0/libexec/python/pyspark/context.py", line 702, in newAPIHadoopRDD  
> > > jconf, batchSize)  
> > > File "/usr/local/Cellar/apache-spark/2.2.0/libexec/python/lib/py4j-0.10.4-src.zip/py4j/java\_gateway.py", line 1133, in **call**  
> > > File "/usr/local/Cellar/apache-spark/2.2.0/libexec/python/pyspark/sql/utils.py", line 63, in deco  
> > > return f(\*a, \*\*kw)  
> > > File "/usr/local/Cellar/apache-spark/2.2.0/libexec/python/lib/py4j-0.10.4-src.zip/py4j/protocol.py", line 319, in get\_return\_value  
> > > py4j.protocol.Py4JJavaError: An error occurred while calling z:org.apache.spark.api.python.PythonRDD.newAPIHadoopRDD.  
> > > : java.lang.ClassCastException: java.lang.Integer cannot be cast to java.lang.String  
> > > at org.apache.spark.api.python.PythonHadoopUtil$$anonfun$mapToConf$1.apply(PythonHadoopUtil.scala:160)  
> > > at org.apache.spark.api.python.PythonHadoopUtil$$anonfun$mapToConf$1.apply(PythonHadoopUtil.scala:160)  
> > > at scala.collection.Iterator$class.foreach(Iterator.scala:893)  
> > > at scala.collection.AbstractIterator.foreach(Iterator.scala:1336)  
> > > at scala.collection.IterableLike$class.foreach(IterableLike.scala:72)  
> > > at scala.collection.AbstractIterable.foreach(Iterable.scala:54)  
> > > at org.apache.spark.api.python.PythonHadoopUtil$.mapToConf(PythonHadoopUtil.scala:160)  
> > > at org.apache.spark.api.python.PythonRDD$.newAPIHadoopRDD(PythonRDD.scala:580)  
> > > at org.apache.spark.api.python.PythonRDD.newAPIHadoopRDD(PythonRDD.scala)  
> > > at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)  
> > > at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)  
> > > at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)  
> > > at java.lang.reflect.Method.invoke(Method.java:498)  
> > > at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)  
> > > at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)  
> > > at py4j.Gateway.invoke(Gateway.java:280)  
> > > at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)  
> > > at py4j.commands.CallCommand.execute(CallCommand.java:79)  
> > > at py4j.GatewayConnection.run(GatewayConnection.java:214)  
> > > at java.lang.Thread.run(Thread.java:745)

Anyone see what I am getting wrong?  
spark version: version 2.2.0  
Elastic Search version: 2.2.1  
No haddop.

Thanks.

---

<div class="post-metadata">

**Author:** ![james.baiera](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/james.baiera/32/10209_2.png) [@james.baiera](https://discuss.elastic.co/u/james.baiera)\
**Post date:** [December 13, 2017, 7:38pm UTC](https://discuss.elastic.co/t/es-hadoop-pyspark-error/110191/2 "2017-12-13T19:38:54Z")

</div>

Could it be that you are specifying the port number as an Integer instead of a String?

```auto
"es.port":9223

```

should be this?

```auto
"es.port":"9223"

```

---

<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:** [January 10, 2018, 7:39pm UTC](https://discuss.elastic.co/t/es-hadoop-pyspark-error/110191/3 "2018-01-10T19:39:00Z")

</div>

This topic was automatically closed 28 days after the last reply. New replies are no longer allowed.
