# Spark structured streaming Elasticsearch integration issue

**URL:** <https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-integration-issue/185366>\
**Category:** Elasticsearch\
**Tags:** es-hadoop\
**Created:** [June 12, 2019, 9:03am UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-integration-issue/185366 "2019-06-12T09:03:11Z")\
**Posts on this page:** 3\
**Page:** 1

<div class="post-metadata">

**Author:** ![Manishvsaraf](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/manishvsaraf/32/47820_2.png) [@Manishvsaraf](https://discuss.elastic.co/u/Manishvsaraf)\
**Post date:** [June 12, 2019, 9:03am UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-integration-issue/185366/1 "2019-06-12T09:03:11Z")

</div>

I am writing a Spark structured streaming application in which data processed with Spark needs be sink'ed to Elasticsearch.  
This is my development environment.  
Hadoop 2.6.0-cdh5.16.1  
Spark version 2.3.0.cloudera4  
elasticsearch 6.8.0

I ran spark-shell as

> spark2-shell --jars /tmp/elasticsearch-hadoop-2.3.2/dist/elasticsearch-hadoop-2.3.2.jar
> 
> import org.apache.spark.SparkContext  
> import org.apache.spark.SparkConf  
> import org.apache.spark.sql.functions.\_  
> import org.apache.spark.sql.SparkSession  
> import org.apache.spark.sql.types.{StructType, StructField, StringType, IntegerType, TimestampType};  
> import java.util.Calendar  
> import org.apache.spark.sql.SparkSession  
> import org.elasticsearch.spark.sql  
> import sys.process.\_
> 
> val checkPointDir = "/tmp/rt/checkpoint/"
> 
> val spark = SparkSession.builder  
> .config("fs.s3n.impl", "org.apache.hadoop.fs.s3native.NativeS3FileSystem")  
> .config("fs.s3n.awsAccessKeyId","aaabbb")  
> .config("fs.s3n.awsSecretAccessKey","aaabbbccc")  
> .config("spark.sql.streaming.checkpointLocation",s"$checkPointDir")  
> .config("es.index.auto.create", "true").getOrCreate()  
> import spark.implicits.\_
> 
> val requestSchema = new StructType().add("log\_type", StringType).add("time\_stamp", StringType).add("host\_name", StringType).add("data\_center", StringType).add("build", StringType).add("ip\_trace", StringType).add("client\_ip", StringType).add("protocol", StringType).add("latency", StringType).add("status", StringType).add("response\_size", StringType).add("request\_id", StringType).add("user\_id", StringType).add("pageview\_id", StringType).add("impression\_id", StringType).add("source\_impression\_id", StringType).add("rnd", StringType).add("publisher\_id", StringType).add("site\_id", StringType).add("zone\_id", StringType).add("slot\_id", StringType).add("tile", StringType).add("content\_id", StringType).add("post\_id", StringType).add("postgroup\_id", StringType).add("brand\_id", StringType).add("provider\_id", StringType).add("geo\_country", StringType).add("geo\_region", StringType).add("geo\_city", StringType).add("geo\_zip\_code", StringType).add("geo\_area\_code", StringType).add("geo\_dma\_code", StringType).add("browser\_group", StringType).add("page\_url", StringType).add("document\_referer", StringType).add("user\_agent", StringType).add("cookies", StringType).add("kvs", StringType).add("notes", StringType).add("request", StringType)  
> val requestDF = spark.readStream.option("delimiter", "\t").format("com.databricks.spark.csv").schema(requestSchema).load("s3n://aa/logs/cc.com/r/year=" + Calendar.getInstance().get(Calendar.YEAR) + "/month=" + "%02d".format(Calendar.getInstance().get(Calendar.MONTH)+1) + "/day=" + "%02d".format(Calendar.getInstance().get(Calendar.DAY\_OF\_MONTH)) + "/hour=" + "%02d".format(Calendar.getInstance().get(Calendar.HOUR\_OF\_DAY)) + "/\*.log")  
> requestDF.writeStream.format("org.elasticsearch.spark.sql").option("es.resource", "rt\_request/doc").option("es.nodes", "localhost").outputMode("Append").start()

I have tried following two ways to sink the data in the DataSet to ES.  
1.ds.writeStream().format("org.elasticsearch.spark.sql").start("spark/orders");  
2.ds.writeStream().format("es").start("rt\_request/doc");  
In both cases I am getting the following error:

Caused by:  
java.lang.UnsupportedOperationException: Data source es does not support streamed writing

java.lang.UnsupportedOperationException: Data source org.elasticsearch.spark.sql does not support streamed writing  
at org.apache.spark.sql.execution.datasources.DataSource.createSink(DataSource.scala:320)  
at org.apache.spark.sql.streaming.DataStreamWriter.start(DataStreamWriter.scala:293)  
... 56 elided

---

<div class="post-metadata">

**Author:** ![Manishvsaraf](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/manishvsaraf/32/47820_2.png) [@Manishvsaraf](https://discuss.elastic.co/u/Manishvsaraf)\
**Post date:** [June 13, 2019, 4:34pm UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-integration-issue/185366/2 "2019-06-13T16:34:29Z")

</div>

ES-hadoop jar version I used is old one elasticsearch-hadoop-2.3.2.jar. we need 6 or above.

Now I use elasticsearch-hadoop-6\* or above jars for it to work as a streaming sink.

I have downloaded it from [https://artifacts.elastic.co/downloads/elasticsearch-hadoop/elasticsearch-hadoop-7.1.1.zip](https://artifacts.elastic.co/downloads/elasticsearch-hadoop/elasticsearch-hadoop-7.1.1.zip)

---

<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 11, 2019, 4:34pm UTC](https://discuss.elastic.co/t/spark-structured-streaming-elasticsearch-integration-issue/185366/3 "2019-07-11T16:34:32Z")

</div>

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