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
