# Spark push-down query enhancement

**URL:** <https://discuss.elastic.co/t/spark-push-down-query-enhancement/63159>\
**Category:** Elasticsearch\
**Tags:** es-hadoop\
**Created:** [October 17, 2016, 7:29am UTC](https://discuss.elastic.co/t/spark-push-down-query-enhancement/63159 "2016-10-17T07:29:09Z")\
**Posts on this page:** 5\
**Page:** 1

<div class="post-metadata">

**Author:** ![aokolnychyi](https://avatars.discourse-cdn.com/v4/letter/a/f05b48/32.png) [@aokolnychyi](https://discuss.elastic.co/u/aokolnychyi)\
**Post date:** [October 17, 2016, 7:29am UTC](https://discuss.elastic.co/t/spark-push-down-query-enhancement/63159/1 "2016-10-17T07:29:09Z")

</div>

Hi,

I am creating a data frame by specifying a filter alongside my index and type name. Everything works fine unless I am trying to additionally filter the resulting data frame on Spark side. According to the documentation, "Note the push down operations apply even when one specifies a query - the connector will enhance it according to the specified SQL". As I understand, my filter should be enhanced but not ignored. However, it seems that my initial query is ignored if the push-down is enabled.

I have the following use case.

Mapping:

```
"mappings": {
  "locationEvent": {
    "properties": {
      "blockId": {
        "type": "string",
        "index": "not_analyzed"
      },
      "deviation": {
        "type": "long"
      },
      "distance": {
        "type": "long"
      },
      "doorsState": {
        "type": "string",
        "index": "not_analyzed"
      },
      "lineId": {
        "type": "string",
        "index": "not_analyzed"
      },
      "routeId": {
        "type": "string",
        "index": "not_analyzed"
      },
      "stopPointId": {
        "type": "string",
        "index": "not_analyzed"
      },
      "time": {
        "type": "date",
        "format": "yyyy-MM-dd'T'HH:mm:ss.SSSZ"
      },
      "tripNumber": {
        "type": "string",
        "index": "not_analyzed"
      }
    }
  }
}

```

Initial filter (later as validEventsQuery):

```
{
  "filter": {
    "bool": {
      "must": [
        {
          "range": {
            "time": {
              "gte": "2016-08-30T00:00:00.000Z",
              "lt": "2016-08-31T02:00:00.000Z"
            }
          }
        },
        {
          "or": [
            {
              "not": {
                "term": {
                  "distance": "0"
                }
              }
            },
            {
              "and": [
            {
                  "term": {
                    "distance": "0"
                  }
                },
                {
                  "term": {
                    "doorsState": "open"
                  }
                }
              ]
            }
          ]
        },
        {
          "exists": {
            "field": "deviation"
          }
        },
        {
          "exists": {
            "field": "distance"
          }
        },
        {
          "exists": {
            "field": "lineId"
          }
        },
        {
          "exists": {
            "field": "tripNumber"
          }
        },
        {
          "exists": {
            "field": "routeId"
          }
        },
        {
          "exists": {
            "field": "doorsState"
          }
        },
        {
          "exists": {
            "field": "time"
          }
        },
        {
          "exists": {
            "field": "stopPointId"
          }
        }
      ]
    }
  }
}

```

1. Default config, no filtering on Spark side.  
`val eventDF = hiveContext.esDF(s"$index/$type", validEventsQuery).count()`  
The result is 480365. ES REST API gives the same number.
2. Default config, additional filtering on Spark side.  
`val eventDF = hiveContext.esDF(s"$index/$type", validEventsQuery).filter('stopPointId !== "").count()`  
The result is 715467. ES REST API returns this number only for the filter stopPointId !== "".

If I disable the push-down, then everything works as expected. I also tried to set "strict" to true and "double.filtering" to true/false but nothing worked for me.

It would be great if anyone can tell what I am doing wrong. Thanks in advance.

Versions:  
sparkVersion = '1.6.1'  
esVersion = '2.3.2'  
esSparkConnectorVersion = '2.3.3'

---

<div class="post-metadata">

**Author:** ![aokolnychyi](https://avatars.discourse-cdn.com/v4/letter/a/f05b48/32.png) [@aokolnychyi](https://discuss.elastic.co/u/aokolnychyi)\
**Post date:** [October 17, 2016, 7:30am UTC](https://discuss.elastic.co/t/spark-push-down-query-enhancement/63159/2 "2016-10-17T07:30:58Z")

</div>

The corresponding enhanced query in the slow search log:

```
    {
      "query": {
        "filtered": {
          "filter": {
            "bool": {
              "must": [
                    {
                  "range": {
                    "time": {
                      "gte": "2016-08-30T00:00:00.000Z",
                      "lt": "2016-08-31T02:00:00.000Z"
                    }
                  }
                },
                {
                  "or": [
                    {
                      "not": {
                        "term": {
                          "distance": "0"
                        }
                      }
                    },
                    {
                      "and": [
                        {
                          "term": {
                            "distance": "0"
                          }
                        },
                        {
                          "term": {
                            "doorsState": "open"
                          }
                        }
                      ]
                    }
                  ]
                },
                {
                  "exists": {
                    "field": "deviation"
                  }
                },
                {
                  "exists": {
                    "field": "distance"
                  }
                },
                {
                  "exists": {
                   "field": "lineId"
                  }
                },
                {
                  "exists": {
                    "field": "tripNumber"
                  }
                },
                {
                  "exists": {
                    "field": "routeId"
                  }
                },
                {
                  "exists": {
                    "field": "doorsState"
                  }
                },
                {
                  "exists": {
                    "field": "time"
                  }
                },
                {
                  "exists": {
                    "field": "stopPointId"
                  }
                }
              ]
            }
          },
          "filter": {
            "and": [
              {
                "not": {
                  "filter": {
                    "query": {
                      "match": {
                        "stopPointId": ""
                      }
                    }
                  }
                }
              }
            ]
          }
        }
      }
    }
```

---

<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:** [October 17, 2016, 8:41pm UTC](https://discuss.elastic.co/t/spark-push-down-query-enhancement/63159/3 "2016-10-17T20:41:41Z")

</div>

It looks like your provided QueryDSL is missing the top level "query" attribute, and so the generated query has two entries for "filter" underneath the "filtered" query that marries the pushdown with the provided query.

If you wrap your initially provided Query DSL in a `{"query":{"filtered":{ ... }}}` block, it should correctly construct the pushdown filter.

---

<div class="post-metadata">

**Author:** ![aokolnychyi](https://avatars.discourse-cdn.com/v4/letter/a/f05b48/32.png) [@aokolnychyi](https://discuss.elastic.co/u/aokolnychyi)\
**Post date:** [October 18, 2016, 2:35pm UTC](https://discuss.elastic.co/t/spark-push-down-query-enhancement/63159/4 "2016-10-18T14:35:12Z")

</div>

Thanks a lot for your reply, I appreciate.  
You are right, it worked.

---

<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/spark-push-down-query-enhancement/63159/5 "2017-07-06T13:22:51Z")

</div>


