# Newbie question about Spark and Elasticsearch

**URL:** <https://discuss.elastic.co/t/newbie-question-about-spark-and-elasticsearch/21146>\
**Category:** Elasticsearch\
**Created:** [December 8, 2014, 2:59pm UTC](https://discuss.elastic.co/t/newbie-question-about-spark-and-elasticsearch/21146 "2014-12-08T14:59:09Z")\
**Posts on this page:** 6\
**Page:** 1

<div class="post-metadata">

**Author:** ![Mohamed\_Lrhazi\_2](https://avatars.discourse-cdn.com/v4/letter/m/d6d6ee/32.png) [@Mohamed\_Lrhazi\_2](https://discuss.elastic.co/u/Mohamed_Lrhazi_2)\
**Post date:** [December 8, 2014, 2:59pm UTC](https://discuss.elastic.co/t/newbie-question-about-spark-and-elasticsearch/21146/1 "2014-12-08T14:59:09Z")

</div>

am trying to understand how spark and ES work... could someone please help  
me answer this question..

val conf = new Configuration()  
conf.set("es.resource", "radio/artists")  
conf.set("es.query", "?q=me\*")  
val esRDD = sc.newHadoopRDD(conf, classOf[EsInputFormat[Text,  
MapWritable]],  
classOf[Text], classOf[MapWritable]))  
val docCount = esRDD.count();

When and where is data being transferred from ES? is it all collected on  
the Spark master node, then partitioned and sent to the worker nodes? or is  
each worker node talking to ES to somehow get a partition of the data?

How does this effectively work?

Thanks a lot,  
Mohamed.

--  
You received this message because you are subscribed to the Google Groups "elasticsearch" group.  
To unsubscribe from this group and stop receiving emails from it, send an email to [elasticsearch+unsubscribe@googlegroups.com](mailto:elasticsearch+unsubscribe@googlegroups.com).  
To view this discussion on the web visit [https://groups.google.com/d/msgid/elasticsearch/CAEU\_gmf9Nt0xn\_0NbzDn\_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com](https://groups.google.com/d/msgid/elasticsearch/CAEU_gmf9Nt0xn_0NbzDn_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com).  
For more options, visit [https://groups.google.com/d/optout](https://groups.google.com/d/optout).

---

<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:** [December 8, 2014, 3:19pm UTC](https://discuss.elastic.co/t/newbie-question-about-spark-and-elasticsearch/21146/2 "2014-12-08T15:19:00Z")

</div>

Hi,

First off I recommend using the native integration (aka the Java/Scala APIs) instead of MapReduce. The latter works but  
the former is better performing and more flexible.

ES works in a similar fashion to the HDFS store - the data doesn't go through the master rather, each task has its own  
partition on works on its own set of data. Behind the scenes we map each worker to an index shard (if there aren't  
enough workers, then some will work across multiple shards).

On 12/8/14 4:59 PM, Mohamed Lrhazi wrote:

> am trying to understand how spark and ES work... could someone please help me answer this question..
> 
> val conf = new Configuration()  
> conf.set("es.resource", "radio/artists")  
> conf.set("es.query", "?q=me\*")  
> val esRDD = sc.newHadoopRDD(conf, classOf[EsInputFormat[Text, MapWritable]],  
> classOf[Text], classOf[MapWritable]))  
> val docCount = esRDD.count();
> 
> When and where is data being transferred from ES? is it all collected on the Spark master node, then partitioned and  
> sent to the worker nodes? or is each worker node talking to ES to somehow get a partition of the data?
> 
> How does this effectively work?
> 
> Thanks a lot,  
> Mohamed.
> 
> --  
> You received this message because you are subscribed to the Google Groups "elasticsearch" group.  
> To unsubscribe from this group and stop receiving emails from it, send an email to  
> [elasticsearch+unsubscribe@googlegroups.com](mailto:elasticsearch+unsubscribe@googlegroups.com) [mailto:elasticsearch+unsubscribe@googlegroups.com](mailto:elasticsearch+unsubscribe@googlegroups.com).  
> To view this discussion on the web visit  
> [https://groups.google.com/d/msgid/elasticsearch/CAEU\_gmf9Nt0xn\_0NbzDn\_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com](https://groups.google.com/d/msgid/elasticsearch/CAEU_gmf9Nt0xn_0NbzDn_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com)  
> [https://groups.google.com/d/msgid/elasticsearch/CAEU\_gmf9Nt0xn\_0NbzDn\_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com?utm\_medium=email&utm\_source=footer](https://groups.google.com/d/msgid/elasticsearch/CAEU_gmf9Nt0xn_0NbzDn_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com?utm_medium=email&utm_source=footer).  
> For more options, visit [https://groups.google.com/d/optout](https://groups.google.com/d/optout).

--  
Costin

--  
You received this message because you are subscribed to the Google Groups "elasticsearch" group.  
To unsubscribe from this group and stop receiving emails from it, send an email to [elasticsearch+unsubscribe@googlegroups.com](mailto:elasticsearch+unsubscribe@googlegroups.com).  
To view this discussion on the web visit [https://groups.google.com/d/msgid/elasticsearch/5485C164.7090405%40gmail.com](https://groups.google.com/d/msgid/elasticsearch/5485C164.7090405%40gmail.com).  
For more options, visit [https://groups.google.com/d/optout](https://groups.google.com/d/optout).

---

<div class="post-metadata">

**Author:** ![Mohamed\_Lrhazi](https://avatars.discourse-cdn.com/v4/letter/m/74df32/32.png) [@Mohamed\_Lrhazi](https://discuss.elastic.co/u/Mohamed_Lrhazi)\
**Post date:** [December 8, 2014, 11:13pm UTC](https://discuss.elastic.co/t/newbie-question-about-spark-and-elasticsearch/21146/3 "2014-12-08T23:13:33Z")

</div>

Great Thanks a lot Costin.

Are people supposed to deploy the Spark workers on the same ES cluster? I  
guess it would make sense for data to remain local and avoid network  
transfers altogether?

Thanks a lot,  
Mohamed.

On Monday, December 8, 2014 10:19:12 AM UTC-5, Costin Leau wrote:

> Hi,
> 
> First off I recommend using the native integration (aka the Java/Scala  
> APIs) instead of MapReduce. The latter works but  
> the former is better performing and more flexible.
> 
> ES works in a similar fashion to the HDFS store - the data doesn't go  
> through the master rather, each task has its own  
> partition on works on its own set of data. Behind the scenes we map each  
> worker to an index shard (if there aren't  
> enough workers, then some will work across multiple shards).
> 
> On 12/8/14 4:59 PM, Mohamed Lrhazi wrote:
> 
> > am trying to understand how spark and ES work... could someone please  
> > help me answer this question..
> > 
> > val conf = new Configuration()  
> > conf.set("es.resource", "radio/artists")  
> > conf.set("es.query", "?q=me\*")  
> > val esRDD = sc.newHadoopRDD(conf, classOf[EsInputFormat[Text,  
> > MapWritable]],  
> > classOf[Text], classOf[MapWritable]))  
> > val docCount = esRDD.count();
> > 
> > When and where is data being transferred from ES? is it all collected on  
> > the Spark master node, then partitioned and  
> > sent to the worker nodes? or is each worker node talking to ES to  
> > somehow get a partition of the data?
> > 
> > How does this effectively work?
> > 
> > Thanks a lot,  
> > Mohamed.
> > 
> > --  
> > You received this message because you are subscribed to the Google  
> > Groups "elasticsearch" group.  
> > To unsubscribe from this group and stop receiving emails from it, send  
> > an email to  
> > [elasticsearc...@googlegroups.com](mailto:elasticsearc...@googlegroups.com) \<javascript:\> \<mailto:  
> > [elasticsearch+unsubscribe@googlegroups.com](mailto:elasticsearch+unsubscribe@googlegroups.com) \<javascript:\>\>.  
> > To view this discussion on the web visit
> 
> [https://groups.google.com/d/msgid/elasticsearch/CAEU\_gmf9Nt0xn\_0NbzDn\_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com](https://groups.google.com/d/msgid/elasticsearch/CAEU_gmf9Nt0xn_0NbzDn_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com)
> 
> > \<  
> > [https://groups.google.com/d/msgid/elasticsearch/CAEU\_gmf9Nt0xn\_0NbzDn\_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com?utm\_medium=email&utm\_source=footer](https://groups.google.com/d/msgid/elasticsearch/CAEU_gmf9Nt0xn_0NbzDn_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com?utm_medium=email&utm_source=footer)\>.
> 
> > For more options, visit [https://groups.google.com/d/optout](https://groups.google.com/d/optout).
> 
> --  
> Costin

--  
You received this message because you are subscribed to the Google Groups "elasticsearch" group.  
To unsubscribe from this group and stop receiving emails from it, send an email to [elasticsearch+unsubscribe@googlegroups.com](mailto:elasticsearch+unsubscribe@googlegroups.com).  
To view this discussion on the web visit [https://groups.google.com/d/msgid/elasticsearch/26361977-b5e1-45fa-b305-e59310e2ce3f%40googlegroups.com](https://groups.google.com/d/msgid/elasticsearch/26361977-b5e1-45fa-b305-e59310e2ce3f%40googlegroups.com).  
For more options, visit [https://groups.google.com/d/optout](https://groups.google.com/d/optout).

---

<div class="post-metadata">

**Author:** ![chrisbtk](https://avatars.discourse-cdn.com/v4/letter/c/a698b9/32.png) [@chrisbtk](https://discuss.elastic.co/u/chrisbtk)\
**Post date:** [December 18, 2014, 6:27pm UTC](https://discuss.elastic.co/t/newbie-question-about-spark-and-elasticsearch/21146/4 "2014-12-18T18:27:42Z")

</div>

Hi,

You recommend the native integration instead of MR and I see on the  
official documentation that MR is recommended to read/write data to ES  
using spark. Spark support Doc  
[http://www.elasticsearch.org/guide/en/elasticsearch/hadoop/2.1.Beta/spark.html](http://www.elasticsearch.org/guide/en/elasticsearch/hadoop/2.1.Beta/spark.html)

what would be the basic piece of code to read data from ES without using MR  
?

I'm currently struggling with EsInputFormat[org.apache.hadoop.io.Text,  
MapWritable] structure.

my code is :

val sc = new SparkContext(...)

val configuration = new Configuration()  
configuration.set("es.nodes", "xxxxxx")  
configuration.set("es.port", "9200")  
configuration.set("es.resource", resource) // my index/type  
configuration.set("es.query", query) //basicaly a match\_all

val esRDD = sc.newAPIHadoopRDD(configuration,  
classOf[EsInputFormat[org.apache.hadoop.io.Text,  
MapWritable]],classOf[org.apache.hadoop.io.Text], classOf[MapWritable])

assume my data is mapped as follow :

{  
"oceannetworks": {  
"mappings": {  
"transcript": {  
"properties": {  
"cruiseID": {  
"type": "string"  
},  
"diveID": {  
"type": "string"  
},  
"filename\_root": {  
"type": "string"  
},  
"id": {  
"type": "string"  
},  
"result": {  
"type": "nested",  
"properties": {  
"begin\_time": {  
"type": "double"  
},  
"confidence": {  
"type": "double"  
},  
"end\_time": {  
"type": "double"  
},  
"location": {  
"type": "geo\_point"  
},  
"word": {  
"type": "string"  
}  
}  
},  
"status": {  
"type": "string"  
},  
"uuid": {  
"type": "string"  
},  
"version": {  
"type": "string"  
}  
}  
}  
}  
}  
}

I'm able to retrieve 1st level information like diveID , cruiseID ... but  
it's not clear how to get the 2nd lvl collection "result". It seams I get a  
WritableArrayWritable but I'm not sure how to handle it.

I get 1st lvl data with these king of code :

val uuids = esRDD.map(\_.\_2.get(new  
org.apache.hadoop.io.Text("uuid")).toString).take(10)

I could use a little bit of help 🙂

thanks.

chris

Le lundi 8 décembre 2014 10:19:12 UTC-5, Costin Leau a écrit :

> Hi,
> 
> First off I recommend using the native integration (aka the Java/Scala  
> APIs) instead of MapReduce. The latter works but  
> the former is better performing and more flexible.
> 
> ES works in a similar fashion to the HDFS store - the data doesn't go  
> through the master rather, each task has its own  
> partition on works on its own set of data. Behind the scenes we map each  
> worker to an index shard (if there aren't  
> enough workers, then some will work across multiple shards).
> 
> On 12/8/14 4:59 PM, Mohamed Lrhazi wrote:
> 
> > am trying to understand how spark and ES work... could someone please  
> > help me answer this question..
> > 
> > val conf = new Configuration()  
> > conf.set("es.resource", "radio/artists")  
> > conf.set("es.query", "?q=me\*")  
> > val esRDD = sc.newHadoopRDD(conf, classOf[EsInputFormat[Text,  
> > MapWritable]],  
> > classOf[Text], classOf[MapWritable]))  
> > val docCount = esRDD.count();
> > 
> > When and where is data being transferred from ES? is it all collected on  
> > the Spark master node, then partitioned and  
> > sent to the worker nodes? or is each worker node talking to ES to  
> > somehow get a partition of the data?
> > 
> > How does this effectively work?
> > 
> > Thanks a lot,  
> > Mohamed.
> > 
> > --  
> > You received this message because you are subscribed to the Google  
> > Groups "elasticsearch" group.  
> > To unsubscribe from this group and stop receiving emails from it, send  
> > an email to  
> > [elasticsearc...@googlegroups.com](mailto:elasticsearc...@googlegroups.com) \<javascript:\> \<mailto:  
> > [elasticsearch+unsubscribe@googlegroups.com](mailto:elasticsearch+unsubscribe@googlegroups.com) \<javascript:\>\>.  
> > To view this discussion on the web visit
> 
> [https://groups.google.com/d/msgid/elasticsearch/CAEU\_gmf9Nt0xn\_0NbzDn\_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com](https://groups.google.com/d/msgid/elasticsearch/CAEU_gmf9Nt0xn_0NbzDn_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com)
> 
> > \<  
> > [https://groups.google.com/d/msgid/elasticsearch/CAEU\_gmf9Nt0xn\_0NbzDn\_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com?utm\_medium=email&utm\_source=footer](https://groups.google.com/d/msgid/elasticsearch/CAEU_gmf9Nt0xn_0NbzDn_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com?utm_medium=email&utm_source=footer)\>.
> 
> > For more options, visit [https://groups.google.com/d/optout](https://groups.google.com/d/optout).
> 
> --  
> Costin

--  
You received this message because you are subscribed to the Google Groups "elasticsearch" group.  
To unsubscribe from this group and stop receiving emails from it, send an email to [elasticsearch+unsubscribe@googlegroups.com](mailto:elasticsearch+unsubscribe@googlegroups.com).  
To view this discussion on the web visit [https://groups.google.com/d/msgid/elasticsearch/959a9e83-5cae-4428-a45d-3ae5266af275%40googlegroups.com](https://groups.google.com/d/msgid/elasticsearch/959a9e83-5cae-4428-a45d-3ae5266af275%40googlegroups.com).  
For more options, visit [https://groups.google.com/d/optout](https://groups.google.com/d/optout).

---

<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:** [December 18, 2014, 8:26pm UTC](https://discuss.elastic.co/t/newbie-question-about-spark-and-elasticsearch/21146/5 "2014-12-18T20:26:58Z")

</div>

> **[Elasticsearch Platform — Find real-time answers at scale](https://www.elastic.co)**
>
> Power insights and outcomes with the Elasticsearch Platform and AI. See into your data and find answers that matter with enterprise solutions designed to help you build, observe, and protect. Try Elasticsearch free today.

On 12/18/14 8:27 PM, chris wrote:

> Hi,
> 
> You recommend the native integration instead of MR and I see on the official documentation that MR is recommended to  
> read/write data to ES using spark. Spark support Doc  
> [http://www.elasticsearch.org/guide/en/elasticsearch/hadoop/2.1.Beta/spark.html](http://www.elasticsearch.org/guide/en/elasticsearch/hadoop/2.1.Beta/spark.html)
> 
> what would be the basic piece of code to read data from ES without using MR ?
> 
> I'm currently struggling with EsInputFormat[org.apache.hadoop.io.Text, MapWritable] structure.
> 
> my code is :
> 
> val sc = new SparkContext(...)
> 
> val configuration = new Configuration()  
> configuration.set("es.nodes", "xxxxxx")  
> configuration.set("es.port", "9200")  
> configuration.set("es.resource", resource) // my index/type  
> configuration.set("es.query", query) //basicaly a match\_all
> 
> val esRDD = sc.newAPIHadoopRDD(configuration, classOf[EsInputFormat[org.apache.hadoop.io.Text,  
> MapWritable]],classOf[org.apache.hadoop.io.Text], classOf[MapWritable])
> 
> assume my data is mapped as follow :
> 
> {  
> "oceannetworks": {  
> "mappings": {  
> "transcript": {  
> "properties": {  
> "cruiseID": {  
> "type": "string"  
> },  
> "diveID": {  
> "type": "string"  
> },  
> "filename\_root": {  
> "type": "string"  
> },  
> "id": {  
> "type": "string"  
> },  
> "result": {  
> "type": "nested",  
> "properties": {  
> "begin\_time": {  
> "type": "double"  
> },  
> "confidence": {  
> "type": "double"  
> },  
> "end\_time": {  
> "type": "double"  
> },  
> "location": {  
> "type": "geo\_point"  
> },  
> "word": {  
> "type": "string"  
> }  
> }  
> },  
> "status": {  
> "type": "string"  
> },  
> "uuid": {  
> "type": "string"  
> },  
> "version": {  
> "type": "string"  
> }  
> }  
> }  
> }  
> }  
> }
> 
> I'm able to retrieve 1st level information like diveID , cruiseID ... but it's not clear how to get the 2nd lvl  
> collection "result". It seams I get a WritableArrayWritable but I'm not sure how to handle it.
> 
> I get 1st lvl data with these king of code :
> 
> val uuids = esRDD.map(\_.\_2.get(new org.apache.hadoop.io.Text("uuid")).toString).take(10)
> 
> I could use a little bit of help 🙂
> 
> thanks.
> 
> chris
> 
> Le lundi 8 décembre 2014 10:19:12 UTC-5, Costin Leau a écrit :
> 
> ```
> Hi,
> 
> First off I recommend using the native integration (aka the Java/Scala APIs) instead of MapReduce. The latter works but
> the former is better performing and more flexible.
> 
> ES works in a similar fashion to the HDFS store - the data doesn't go through the master rather, each task has its own
> partition on works on its own set of data. Behind the scenes we map each worker to an index shard (if there aren't
> enough workers, then some will work across multiple shards).
> 
> On 12/8/14 4:59 PM, Mohamed Lrhazi wrote:
> > am trying to understand how spark and ES work... could someone please help me answer this question..
> >
> > val conf = new Configuration()
> > conf.set("es.resource", "radio/artists")
> > conf.set("es.query", "?q=me*")
> > val esRDD = sc.newHadoopRDD(conf, classOf[EsInputFormat[Text, MapWritable]],
> > classOf[Text], classOf[MapWritable]))
> > val docCount = esRDD.count();
> >
> >
> > When and where is data being transferred from ES? is it all collected on the Spark master node, then partitioned and
> > sent to the worker nodes? or is each worker node talking to ES to somehow get a partition of the data?
> >
> > How does this effectively work?
> >
> > Thanks a lot,
> > Mohamed.
> >
> > --
> > You received this message because you are subscribed to the Google Groups "elasticsearch" group.
> > To unsubscribe from this group and stop receiving emails from it, send an email to
> >elasticsearc...@googlegroups.com <javascript:> <mailto:elasticsearch+unsubscribe@googlegroups.com <javascript:>>.
> > To view this discussion on the web visit
> >https://groups.google.com/d/msgid/elasticsearch/CAEU_gmf9Nt0xn_0NbzDn_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com
> <https://groups.google.com/d/msgid/elasticsearch/CAEU_gmf9Nt0xn_0NbzDn_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com>
> > <https://groups.google.com/d/msgid/elasticsearch/CAEU_gmf9Nt0xn_0NbzDn_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com?utm_medium=email&utm_source=footer
> <https://groups.google.com/d/msgid/elasticsearch/CAEU_gmf9Nt0xn_0NbzDn_moRWUT96uWYf4cicJdZik3r0Zz8XA%40mail.gmail.com?utm_medium=email&utm_source=footer>>.
> 
> > For more options, visithttps://groups.google.com/d/optout <https://groups.google.com/d/optout>.
> 
> --
> Costin
> 
> ```
> 
> --  
> You received this message because you are subscribed to the Google Groups "elasticsearch" group.  
> To unsubscribe from this group and stop receiving emails from it, send an email to  
> [elasticsearch+unsubscribe@googlegroups.com](mailto:elasticsearch+unsubscribe@googlegroups.com) [mailto:elasticsearch+unsubscribe@googlegroups.com](mailto:elasticsearch+unsubscribe@googlegroups.com).  
> To view this discussion on the web visit  
> [https://groups.google.com/d/msgid/elasticsearch/959a9e83-5cae-4428-a45d-3ae5266af275%40googlegroups.com](https://groups.google.com/d/msgid/elasticsearch/959a9e83-5cae-4428-a45d-3ae5266af275%40googlegroups.com)  
> [https://groups.google.com/d/msgid/elasticsearch/959a9e83-5cae-4428-a45d-3ae5266af275%40googlegroups.com?utm\_medium=email&utm\_source=footer](https://groups.google.com/d/msgid/elasticsearch/959a9e83-5cae-4428-a45d-3ae5266af275%40googlegroups.com?utm_medium=email&utm_source=footer).  
> For more options, visit [https://groups.google.com/d/optout](https://groups.google.com/d/optout).

--  
Costin

--  
You received this message because you are subscribed to the Google Groups "elasticsearch" group.  
To unsubscribe from this group and stop receiving emails from it, send an email to [elasticsearch+unsubscribe@googlegroups.com](mailto:elasticsearch+unsubscribe@googlegroups.com).  
To view this discussion on the web visit [https://groups.google.com/d/msgid/elasticsearch/54933892.9040703%40gmail.com](https://groups.google.com/d/msgid/elasticsearch/54933892.9040703%40gmail.com).  
For more options, visit [https://groups.google.com/d/optout](https://groups.google.com/d/optout).

---

<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, 12:42am UTC](https://discuss.elastic.co/t/newbie-question-about-spark-and-elasticsearch/21146/6 "2017-07-06T00:42:48Z")

</div>


