# Data ingestion into ElasticSearch from Spark Structured Streaming

**URL:** <https://discuss.elastic.co/t/data-ingestion-into-elasticsearch-from-spark-structured-streaming/196075>\
**Category:** Elasticsearch\
**Tags:** es-hadoop\
**Created:** [August 21, 2019, 9:28am UTC](https://discuss.elastic.co/t/data-ingestion-into-elasticsearch-from-spark-structured-streaming/196075 "2019-08-21T09:28:17Z")\
**Posts on this page:** 4\
**Page:** 1

<div class="post-metadata">

**Author:** ![SandeepReddy](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/sandeepreddy/32/52682_2.png) [@SandeepReddy](https://discuss.elastic.co/u/SandeepReddy)\
**Post date:** [August 21, 2019, 9:28am UTC](https://discuss.elastic.co/t/data-ingestion-into-elasticsearch-from-spark-structured-streaming/196075/1 "2019-08-21T09:28:17Z")

</div>

Hi,

I have created dataset/dataframe using the watermark and window function and writing the output to ElasticSearch is not working.

However dataset/dataframe created without watermark and window inserts data into ElasticSearch.

Please find the code snippet.

val df = dsLog1.withWatermark("time","3 minutes").groupBy(window(col("time"),"3 minutes","1 minute"))  
.agg(count(col("column\_name")))

df.writeStream  
.outputMode("append")  
.format("org.elasticsearch.spark.sql")  
.option("es.nodes", "localhost")  
.option("es.port", "9200")  
//.option("checkpointLocation", "/tmp")  
.option("es.resource","test4/\_doc")  
.start()

Thanks in Advance for the help.

---

<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:** [August 23, 2019, 4:39pm UTC](https://discuss.elastic.co/t/data-ingestion-into-elasticsearch-from-spark-structured-streaming/196075/2 "2019-08-23T16:39:26Z")

</div>

Welcome to the forums!

> I have created dataset/dataframe using the watermark and window function and writing the output to Elasticsearch is not working.

Can you include what is not working about the job?

---

<div class="post-metadata">

**Author:** ![SandeepReddy](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/sandeepreddy/32/52682_2.png) [@SandeepReddy](https://discuss.elastic.co/u/SandeepReddy)\
**Post date:** [August 24, 2019, 4:36pm UTC](https://discuss.elastic.co/t/data-ingestion-into-elasticsearch-from-spark-structured-streaming/196075/3 "2019-08-24T16:36:44Z")

</div>

Hi James,

Thanks for your reply. Please find the details as below:

I am using Spark v2.4 and ElasticSearch v6.8 and my index auto creation is "true"

#######Code Snippet that is working #####################

val dsIn = spark.readStream  
.format(KAFKA\_FORMAT)  
.option(KAFKA\_HOSTS\_PROP, Configuration.getStringList(KAFKA\_HOSTS\_CONF).toArray.mkString(","))  
.option(CONSUMER\_SUBSCRIBE\_PROP, Configuration.getString(CONSUMER\_SUBSCRIBE\_CONF))  
.option(CONSUMER\_GROUPID\_PROP, Configuration.getString(CONSUMER\_GROUPID\_CONF))  
.option(CONSUMER\_ENABLE\_AUTOCOMMIT\_PROP, Configuration.getBoolean(CONSUMER\_ENABLE\_AUTOCOMMIT\_CONF))  
.option(CONSUMER\_STARTING\_OFFSET\_PROP, Configuration.getString(CONSUMER\_STARTING\_OFFSET\_CONF))  
.option(CONSUMER\_FAILONLOSS\_PROP, false)  
.load()  
.selectExpr("CAST(value AS STRING)")  
.as[(String)]  
.select(from\_json($"value", inputSchema).as("value"))

dsIn.writeStream  
.outputMode("append")  
.format("org.elasticsearch.spark.sql")  
.option("es.nodes", "localhost")  
.option("es.port", "9200")  
.option("checkpointLocation", "/tmp")  
// .option("es.mapping.id", "id")  
//.option("es.resource","alta\_alerts\_test/testdata")  
.start("alta\_alerts\_dummy\_2/\_doc")  
// .awaitTermination()

######### Code Snippet that is not working ###############

val dsIn = spark.readStream  
.format(KAFKA\_FORMAT)  
.option(KAFKA\_HOSTS\_PROP, Configuration.getStringList(KAFKA\_HOSTS\_CONF).toArray.mkString(","))  
.option(CONSUMER\_SUBSCRIBE\_PROP, Configuration.getString(CONSUMER\_SUBSCRIBE\_CONF))  
.option(CONSUMER\_GROUPID\_PROP, Configuration.getString(CONSUMER\_GROUPID\_CONF))  
.option(CONSUMER\_ENABLE\_AUTOCOMMIT\_PROP, Configuration.getBoolean(CONSUMER\_ENABLE\_AUTOCOMMIT\_CONF))  
.option(CONSUMER\_STARTING\_OFFSET\_PROP, Configuration.getString(CONSUMER\_STARTING\_OFFSET\_CONF))  
.option(CONSUMER\_FAILONLOSS\_PROP, false)  
.load()  
.selectExpr("CAST(value AS STRING)")  
.as[(String)]  
.select(from\_json($"value", inputSchema).as("value"))

## Group By few columns with the specific time window

val requestCounts = dsIn  
.withWatermark("receivedtime", "10 minutes")  
.groupBy("a", "b", "c", "d", "e", "f", "g", "h", "i", "j" as ("application"), window($"receivedtime", "10 minutes", "5 minutes") as "window").count()

requestCounts.writeStream  
.outputMode("append")  
.format("org.elasticsearch.spark.sql")  
.option("es.nodes", "localhost")  
.option("es.port", "9200")  
.option("checkpointLocation", "/tmp")  
// .option("es.mapping.id", "id")  
//.option("es.resource","alta\_alerts\_test/testdata")  
.start("alta\_alerts\_dummy\_2/\_doc")  
// .awaitTermination()

Thanks in advance for your help.

Regards,  
Sandeep Reddy

---

<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:** [September 21, 2019, 4:36pm UTC](https://discuss.elastic.co/t/data-ingestion-into-elasticsearch-from-spark-structured-streaming/196075/4 "2019-09-21T16:36:48Z")

</div>

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