# Spark sql filter pushdown

**URL:** <https://discuss.elastic.co/t/spark-sql-filter-pushdown/46985>\
**Category:** Elasticsearch\
**Tags:** es-hadoop\
**Created:** [April 11, 2016, 11:59am UTC](https://discuss.elastic.co/t/spark-sql-filter-pushdown/46985 "2016-04-11T11:59:26Z")\
**Posts on this page:** 7\
**Page:** 1

<div class="post-metadata">

**Author:** ![vincent\_gromakowski](https://avatars.discourse-cdn.com/v4/letter/v/59ef9b/32.png) [@vincent\_gromakowski](https://discuss.elastic.co/u/vincent_gromakowski)\
**Post date:** [April 11, 2016, 11:59am UTC](https://discuss.elastic.co/t/spark-sql-filter-pushdown/46985/1 "2016-04-11T11:59:26Z")

</div>

I cannot make the filter to be pushdown to ES when reading from a 1 billion docs index.  
Here is the code:

val options = Map("pushdown" -\> "true", "double.filtering" -\> "false")  
val es=sqlContext.read.format("es").options(options).load("myIndex/myDocs").filter($"id".equalTo("idvalue"))

I can't see anything log that indicates a pushdown:  
16/04/11 13:54:42 DEBUG ScalaEsRowRDD: Using pre-defined reader serializer [org.elasticsearch.spark.sql.ScalaRowValueReader] as default  
16/04/11 13:54:42 DEBUG ScalaEsRowRDD: Using pre-defined reader serializer [org.elasticsearch.spark.sql.ScalaRowValueReader] as default  
16/04/11 13:54:42 DEBUG ScalaEsRowRDD: Partition reader instance [EsPartition [node=[wZ96KouxSueLH7petpE1Gw/opvaames04|xxxxxxxxxx:9200],shard=1]] assigned to [wZ96KouxSueLH7petpE1Gw]:[9200]  
16/04/11 13:54:42 DEBUG ScalaEsRowRDD: Partition reader instance [EsPartition [node=[-Edj1O1LTcK6\_cz1XjOEGA/opvaames06|xxxxxxxxxxx:9200],shard=5]] assigned to [-Edj1O1LTcK6\_cz1XjOEGA]:[9200]

I have also tried this way, it doesn't return any doc but it's quick  
val es= sqlContext.esDF("myIndex/myDocs","?q=idValue")

My conf is Spark 1.6.0 with elasticsearch connector 2.2.0 and elasticsearch 2.1.1

Do you have any advise or method to test simple pushdown ?

---

<div class="post-metadata">

**Author:** ![vincent\_gromakowski](https://avatars.discourse-cdn.com/v4/letter/v/59ef9b/32.png) [@vincent\_gromakowski](https://discuss.elastic.co/u/vincent_gromakowski)\
**Post date:** [April 11, 2016, 1:38pm UTC](https://discuss.elastic.co/t/spark-sql-filter-pushdown/46985/2 "2016-04-11T13:38:15Z")

</div>

I have done some tests and I don't get into this issue while filtering on a flat dataframe structure. My problem is that I am building my documents into ES joining multiple datasets. My structure is like this :  
myObject  
|----- mySubObject1  
|--------------------- field1  
|--------------------- field2  
|----- mySubObject2  
|--------------------- field3

so I need to do :  
val test = sqlContext.read.format("es").load("myIndex/myDocs").filter($"mySubOject1.field1".equalTo("value"))

The pushdown doesn't occur...

---

<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 11, 2016, 2:36pm UTC](https://discuss.elastic.co/t/spark-sql-filter-pushdown/46985/3 "2016-04-11T14:36:41Z")

</div>

First off, try using ES-Hadoop 2.3.0 or 2.2.1.  
Regarding the pushdown, you can enable logging on the spark package (looks like you did) on TRACE level and see whether something shows up. Try to do a simple test (`select A from X where B > O`) and you should see the query being translated.  
Note that pushdown happens if Spark triggers it. I'm not clear what you mean by building "documents into ES joining".

---

<div class="post-metadata">

**Author:** ![vincent\_gromakowski](https://avatars.discourse-cdn.com/v4/letter/v/59ef9b/32.png) [@vincent\_gromakowski](https://discuss.elastic.co/u/vincent_gromakowski)\
**Post date:** [April 11, 2016, 3:13pm UTC](https://discuss.elastic.co/t/spark-sql-filter-pushdown/46985/4 "2016-04-11T15:13:50Z")

</div>

I have already upgraded the connector in 2.3 without better results. I said I am building my documents joining 2 documents : I build myObject with a join between mySubOject1 and mySubObject2 so that's why my data structure isn't flat.  
I will have to flatten my schema which requires an additional processing step in my spark job...

---

<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 11, 2016, 3:47pm UTC](https://discuss.elastic.co/t/spark-sql-filter-pushdown/46985/5 "2016-04-11T15:47:32Z")

</div>

The pushdown applies only to documents that are stored in ES. If your documents are joined (which by the way, is an operation not pushed by Spark) it means they exist in Spark, hence there's nothing really to be pushed down.

---

<div class="post-metadata">

**Author:** ![vincent\_gromakowski](https://avatars.discourse-cdn.com/v4/letter/v/59ef9b/32.png) [@vincent\_gromakowski](https://discuss.elastic.co/u/vincent_gromakowski)\
**Post date:** [April 11, 2016, 4:28pm UTC](https://discuss.elastic.co/t/spark-sql-filter-pushdown/46985/6 "2016-04-11T16:28:50Z")

</div>

Of course my docs are in ES... Anyway the problem is with non flat data structure : sub documents in documents (object.subObject in Spark SQL), the filter isn't pushdown.

---

<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:25pm UTC](https://discuss.elastic.co/t/spark-sql-filter-pushdown/46985/7 "2017-07-06T13:25:13Z")

</div>


