# Spark: optimize partitioning

**URL:** <https://discuss.elastic.co/t/spark-optimize-partitioning/107425>\
**Category:** Elasticsearch\
**Tags:** es-hadoop\
**Created:** [November 13, 2017, 4:01pm UTC](https://discuss.elastic.co/t/spark-optimize-partitioning/107425 "2017-11-13T16:01:55Z")\
**Posts on this page:** 4\
**Page:** 1

<div class="post-metadata">

**Author:** ![amyantis](https://avatars.discourse-cdn.com/v4/letter/a/b9e5f3/32.png) [@amyantis](https://discuss.elastic.co/u/amyantis)\
**Post date:** [November 13, 2017, 4:01pm UTC](https://discuss.elastic.co/t/spark-optimize-partitioning/107425/1 "2017-11-13T16:01:55Z")

</div>

Hi,  
In Spark, I am reading data from Elasticsearch. I need to group documents over the value of a specific field in order to map those groups to a function. My mapper is therefore waiting for a list of documents that as one of the fields identical for each document.

Here is an example of what I am trying to optimize.

I have the following documents in an index:

```json
[{
    "user" : "kimchy",
    "message" : "trying out Elasticsearch"
},
{
    "user" : "kimchy",
    "message" : "hello world"
},
{
    "user" : "michelle",
    "message" : "Hi world"
}

```

I used `user` as routing so that all messages of a user are in a same shard.

In spark, I am reading this index and I wan't to partition messages per user (and not per shard only).

My mapper function is something like

```python
def mapper(messages, user):
    # If user is none, get it from first message
    # 1. compute some data with those messages
    # 2. save some data to ES

```

For now, this is what I am doing in `pyspark`:

```python
df = sqlContext.read.format("org.elasticsearch.spark.sql").load("tweet/message")
rdd = df.rdd.groupBy(lambda x: x.user)
rdd.mapPartitions(mapper, preservesPartitioning=True).reduce(reducer)

```

Is there anything I can do to optimize the partitioning?

---

<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 4, 2017, 8:06pm UTC](https://discuss.elastic.co/t/spark-optimize-partitioning/107425/2 "2017-12-04T20:06:58Z")

</div>

As it stands we do not support defining partitions by a field value in the connector. It's unlikely that this will be supported going forward as it's a very niche feature. From what you are describing though, it seems that this is a decent approach for making sure that all user information is available per partition. Adding parent-child, nested or join fields will probably make the performance worse in terms of querying the data, so I'd advise against that.

---

<div class="post-metadata">

**Author:** ![amyantis](https://avatars.discourse-cdn.com/v4/letter/a/b9e5f3/32.png) [@amyantis](https://discuss.elastic.co/u/amyantis)\
**Post date:** [December 5, 2017, 8:28am UTC](https://discuss.elastic.co/t/spark-optimize-partitioning/107425/3 "2017-12-05T08:28:41Z")

</div>

Thank you for your answer.

---

<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 2, 2018, 8:28am UTC](https://discuss.elastic.co/t/spark-optimize-partitioning/107425/4 "2018-01-02T08:28:57Z")

</div>

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