您好,登錄后才能下訂單哦!
一:先寫map類
import sys for line in sys.stdin: line = line.strip( ) words = line.split( ) for word in words: print('%s\t%s' % (word, 1))
二:寫reduce類
import sys current_word = None current_count = 0 word = None for line in sys.stdin: line = line.strip() word, count = line.split('\t',1) try: count = int(count) except ValueError: continue if current_word == word: current_count += count else: if current_word: print('%s\t%s' % (current_word,current_count)) current_count = count current_word = word if current_word == word: print('%s\t%s' % (current_word,current_count))
三:利用hadoop Streaming執行Python的內容。
hadoop jar /home/hadoop/hadoop-2.6.0-cdh6.5.2/share/hadoop/tools/lib/hadoop-streaming-2.6.0-cdh6.5.2.jar -input /user/hadoop/aa.txt -output /user/hadoop/python_output -mapper "python mapper.py" -reducer "python reducer.py" -file mapper.py -file reducer.py
說明:
輸入和輸出路徑,本身就是hdfs上的,不需要特殊指定hdfs。
不加×××部分的引號的話,會報錯誤:
Error: java.lang.RuntimeException: PipeMapRed.waitOutputThreads(): subprocess failed with code 2
不加粉色部分的內容的話,會報錯誤:
Error: java.lang.RuntimeException: Error in configuring object
免責聲明:本站發布的內容(圖片、視頻和文字)以原創、轉載和分享為主,文章觀點不代表本網站立場,如果涉及侵權請聯系站長郵箱:is@yisu.com進行舉報,并提供相關證據,一經查實,將立刻刪除涉嫌侵權內容。