# \[Spark Structured streaming\] Elasticsearch sink index name

**URL:** <https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-sink-index-name/152449>\
**Category:** Elasticsearch\
**Tags:** es-hadoop\
**Created:** [October 15, 2018, 8:23am UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-sink-index-name/152449 "2018-10-15T08:23:27Z")\
**Posts on this page:** 9\
**Page:** 1

<div class="post-metadata">

**Author:** ![VincentL](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/vincentl/32/36514_2.png) [@VincentL](https://discuss.elastic.co/u/VincentL)\
**Post date:** [October 15, 2018, 8:23am UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-sink-index-name/152449/1 "2018-10-15T08:23:27Z")

</div>

Hello,

I'm currently working with Spark Structured streaming, reading data from Kafka, and writing the output to an Elasticsearch sink. Everything is working fine, except that I would like the index name to be generated based on the current date/time, as it can be done in Logstash for instance.  
Is there any way to do it ? Below is my code for reference (pySpark), latest try 🙂 :

```
def generateIndexName():
        return("es_spark_" + datetime.now().strftime("%Y_%m_%d__%H_%M"))
    
query = kafka_stream.writeStream \
    .outputMode("append") \
    .queryName("writing_to_es") \
    .format("org.elasticsearch.spark.sql") \
    .option("checkpointLocation", "C:/TEMP/") \
    .option("es.resource", generateIndexName() + "/es_spark") \
    .option("es.nodes", "192.168.1.1:9200") \
    .start()
query.awaitTermination()

```

The index name is generated when the code is first executed, and then the index name is never modified.

Thanks in advance if you have any idea how to figure this out.

Regards.

Vincent.

---

<div class="post-metadata">

**Author:** ![Christian\_Dahlqvist](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/christian_dahlqvist/32/4617_2.png) [@Christian\_Dahlqvist](https://discuss.elastic.co/u/Christian_Dahlqvist)\
**Post date:** [October 15, 2018, 9:07am UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-sink-index-name/152449/2 "2018-10-15T09:07:01Z")

</div>

In Logstash the index name is determined based on the `@timestamp` field for every event, which does not seem to be what you are doing here. I do not know how to do that in your case, so will leave that for someone else to comment on, but would strongly advice against creating an index per minute as that is likely to cause a lot of problems.

---

<div class="post-metadata">

**Author:** ![VincentL](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/vincentl/32/36514_2.png) [@VincentL](https://discuss.elastic.co/u/VincentL)\
**Post date:** [October 15, 2018, 9:51am UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-sink-index-name/152449/3 "2018-10-15T09:51:42Z")

</div>

Thanks for your reply. You're right, using the timestamp of the event could be a better idea than using the time the event is processed by spark.

Actually the "one index per minute" was configured only for testing purpose to check the frequency of index creation.

---

<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 23, 2018, 3:08pm UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-sink-index-name/152449/4 "2018-10-23T15:08:45Z")

</div>

As christian has mentioned above, ES-Hadoop will not create an index name using the current time, but rather will let you give a field name on your document to use in the index name creation. This field can contain a timestamp, which is usually what users want when they are using time based indices.

---

<div class="post-metadata">

**Author:** ![VincentL](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/vincentl/32/36514_2.png) [@VincentL](https://discuss.elastic.co/u/VincentL)\
**Post date:** [October 29, 2018, 9:36am UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-sink-index-name/152449/5 "2018-10-29T09:36:20Z")

</div>

Thank you James.  
My concern is about the index size since I have to store several hundreds of gigabytes per day. I don't want to restart my Spark Streaming job every hour or so to modify the index name.  
Still, as you said that ES-Hadoop can not dynamicaly generate the index name, I assume there is no other workaround.

---

<div class="post-metadata">

**Author:** ![Christian\_Dahlqvist](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/christian_dahlqvist/32/4617_2.png) [@Christian\_Dahlqvist](https://discuss.elastic.co/u/Christian_Dahlqvist)\
**Post date:** [October 29, 2018, 9:43am UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-sink-index-name/152449/6 "2018-10-29T09:43:26Z")

</div>

Might it be possible to [create a rollover index](https://www.elastic.co/guide/en/elasticsearch/reference/6.4/indices-rollover-index.html) and index into this? This would give you a fixed alias to index into and allow the changing of indices to occur in the background based on size and/or age.

---

<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 29, 2018, 2:16pm UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-sink-index-name/152449/7 "2018-10-29T14:16:57Z")

</div>

@VincentL There shouldn't be anything to stop you from adding a final step to your spark job that adds the current time to a field on your document, and then using that field in your index name pattern for ES-Hadoop. If you don't want that time field to end up in ES, you can tell ES-Hadoop to not include it in the final document by setting `es.mapping.exclude` to that field name.

---

<div class="post-metadata">

**Author:** ![VincentL](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/vincentl/32/36514_2.png) [@VincentL](https://discuss.elastic.co/u/VincentL)\
**Post date:** [October 31, 2018, 8:42am UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-sink-index-name/152449/8 "2018-10-31T08:42:48Z")

</div>

Thank you very much @james.baiera , I searched the documentation to figure out how to use a field this way, and eventually found the answer here : [https://www.elastic.co/guide/en/elasticsearch/hadoop/current/spark.html#spark-streaming-write-dyn](https://www.elastic.co/guide/en/elasticsearch/hadoop/current/spark.html#spark-streaming-write-dyn) .  
I already have several fields in my documents providing timestamps, I just had to tweak one a little bit so it can be properly parsed.  
@Christian_Dahlqvist : thanks for the rollover index, which might do the job as well. Still, I wanted to include this mecanic in spark, for no valid reason 🙂 .

---

<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:** [December 19, 2018, 10:48am UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-sink-index-name/152449/10 "2018-12-19T10:48:46Z")

</div>

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