Smalyshev has uploaded a new change for review.

  https://gerrit.wikimedia.org/r/248564

Change subject: Add first sketch of the transfer script
......................................................................

Add first sketch of the transfer script

Run with:
spark-submit --master yarn transferToES.py -s /tmp/source.parquet -u 
http://my.elastic/wikiname/index
Bug: T113440

Change-Id: Id4df5bffca445b396c875329bce005327110ddab
---
A transferToES.py
1 file changed, 58 insertions(+), 0 deletions(-)


  git pull ssh://gerrit.wikimedia.org:29418/wikimedia/discovery/analytics 
refs/changes/64/248564/1

diff --git a/transferToES.py b/transferToES.py
new file mode 100644
index 0000000..8cfcfdc
--- /dev/null
+++ b/transferToES.py
@@ -0,0 +1,58 @@
+# A very primitive script to transfer from hive data files
+# to ElasticSearch HTTP endpoint
+from pyspark import SparkContext
+from pyspark.sql import SQLContext, Row
+from optparse import OptionParser
+import requests
+import json
+import logging
+
+oparser = OptionParser()
+oparser.add_option("-s", "--source", dest="source", help="source for the 
data", metavar="SOURCE")
+oparser.add_option("-u", "--url", dest="url", help="URL to send the data to", 
metavar="URL")
+oparser.add_option("-p", "--part", dest="part", help="Number of partitions", 
metavar="NUM", default=2)
+(options, args) = oparser.parse_args()
+
+SOURCE = options.source
+TARGET = options.url
+PARTITIONS = options.part
+print "Transferring from %s to %s, %d partitions" %(SOURCE, TARGET, PARTITIONS)
+
+if __name__ == "__main__":
+    sc = SparkContext(appName="Send To ES")
+
+    sqlContext = SQLContext(sc)
+
+    def documentData(document):
+        """
+        Create textual representation of the document data for one document
+        """
+        updateData = {"update": {"_id": document.id}}
+        docDict = document.asDict()
+        del docDict['id']
+        updateDoc = {"doc":  docDict}
+        return json.dumps(updateData) + "\n" + json.dumps(updateDoc) + "\n"
+        
+    def sendDataToES(data):
+        """
+        Send data to ES server
+        """
+        #print data
+        #logging.info(data)
+        r = requests.put(TARGET, data=data)
+        print r.status_code
+        r.close()
+        # FIXME: add error handling
+        
+    def sendDocumentsToES(documents):
+        """
+        Send a set of documents to ES
+        """
+        data = ""
+        for document in documents:
+            data += documentData(document)
+        sendDataToES(data)
+
+    data = sqlContext.load(SOURCE)
+    print "Count: %d\n" % data.count()
+    data.repartition(PARTITIONS).foreachPartition(sendDocumentsToES)

-- 
To view, visit https://gerrit.wikimedia.org/r/248564
To unsubscribe, visit https://gerrit.wikimedia.org/r/settings

Gerrit-MessageType: newchange
Gerrit-Change-Id: Id4df5bffca445b396c875329bce005327110ddab
Gerrit-PatchSet: 1
Gerrit-Project: wikimedia/discovery/analytics
Gerrit-Branch: master
Gerrit-Owner: Smalyshev <[email protected]>

_______________________________________________
MediaWiki-commits mailing list
[email protected]
https://lists.wikimedia.org/mailman/listinfo/mediawiki-commits

Reply via email to