# Pyspark write to Elasticsearch from Kafka

**URL:** <https://discuss.elastic.co/t/pyspark-write-to-elasticsearch-from-kafka/79877>\
**Category:** Elasticsearch\
**Tags:** es-hadoop\
**Created:** [March 24, 2017, 11:00am UTC](https://discuss.elastic.co/t/pyspark-write-to-elasticsearch-from-kafka/79877 "2017-03-24T11:00:03Z")\
**Posts on this page:** 3\
**Page:** 1

<div class="post-metadata">

**Author:** ![newbie\_here](https://avatars.discourse-cdn.com/v4/letter/n/dbc845/32.png) [@newbie\_here](https://discuss.elastic.co/u/newbie_here)\
**Post date:** [March 24, 2017, 11:00am UTC](https://discuss.elastic.co/t/pyspark-write-to-elasticsearch-from-kafka/79877/1 "2017-03-24T11:00:03Z")

</div>

Hi,

I am using the following code in pyspakr to write data into Elasticsearch from Kafka

import pyspark  
from pyspark.sql import SQLContext  
from pyspark import SparkContext, SparkConf  
from pyspark.streaming import StreamingContext  
from pyspark.streaming.kafka import KafkaUtils  
from pyspark.sql.types import StructType, StructField, StringType, IntegerType  
import json  
from kafka import KafkaConsumer, KafkaClient  
from pyspark.sql.functions import explode, split  
import pandas as pd  
from collections import OrderedDict  
from datetime import date

conf = pyspark.SparkConf()  
conf.setMaster('mesos://172.20.1.157:5050')  
sc = pyspark.SparkContext(conf=conf)  
sqlContext = SQLContext(sc)  
ssc = StreamingContext(sc, 2)  
kafkaParams = {'metadata.broker.list': '172.20.1.163:9092', 'auto.offset.reset': 'smallest'} # kafka parameters  
topics = ['json\_topic']  
kafka\_stream = KafkaUtils.createDirectStream(ssc,topics,kafkaParams)  
parsed = kafka\_stream.map(lambda (k,v) : json.loads(v))  
def function(x):  
y = x.collect()  
for d in y :  
print d  
rdd\_json = json.dumps(d)  
print rdd\_json  
rdd\_json.write.format('org.elasticsearch.spark.sql').mode('append').option('es.index.auto.create','true').option('es.resource', 'sql6/swati').save()

parsed.foreachRDD(lambda x :function(x))  
ssc.start()  
ssc.awaittermination()  
ssc.stop()

The output of rdd\_json comes as below :-  
{"Enrolment\_Date": "2008-01-01", "Freq": 78, "Group": "Recorded Data"}

But I am not able to write it my elasticsearch cluster using the following code :-

rdd\_json.write.format('org.elasticsearch.spark.sql').mode('append').option('es.index.auto.create','true').option('es.resource', 'sql6/type').save()

Any ideas as to how to write it into elasticsearch or is any change in the code is required?

Thanks in advance

---

<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:** [March 27, 2017, 3:53pm UTC](https://discuss.elastic.co/t/pyspark-write-to-elasticsearch-from-kafka/79877/2 "2017-03-27T15:53:35Z")

</div>

@newbie_here Could you include some sort of error trace that you are seeing, or does the command fail with no feedback?

---

<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:** [April 24, 2017, 3:53pm UTC](https://discuss.elastic.co/t/pyspark-write-to-elasticsearch-from-kafka/79877/3 "2017-04-24T15:53:37Z")

</div>

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