Pokazywanie postów oznaczonych etykietą multiprocessing. Pokaż wszystkie posty
Pokazywanie postów oznaczonych etykietą multiprocessing. Pokaż wszystkie posty

2017/01/23

Jak działa algorytm map-reduce?

Map, Reduce i wielowątkowość w Pythonie

MapReduce jest własnościową platformą Google do przetwarzania równoległego. Nazwa sugeruje związek z operacjami map i reduce, jednak przeglądając przykłady, spotkałem się od razu z pojęciami mapperów, reduktorów, partycjonerów itd, bezpośrednio powiązanych z implementacjami realizującymi algorytm (jak np. hadoopowy mapreduce).
Poniższy przykład ilustruje podstawy algorytmu przy pomocy elementarnych elementów biblioteki standardowej Pythona - metody map z multiprocessing.pool i reduce z functools (w Pythonie 2 będącej funkcją podstawową).
Program wieloprocesowo zlicza wystąpienia poszczególnych znaków w zadanym ciągu liter.
'''prosty wielowatkowy program
ilustrujacy idee algorytmu map-reduce'''

import multiprocessing
from functools import reduce

def policz_litery(lancuch):
    '''funkcja map'''
    unikalne_znaki = set(lancuch)
    wystapienia = { k:0 for k in unikalne_znaki }
    for litera in lancuch:
        wystapienia[litera] += 1
    return wystapienia

def polacz_slowniki(slownik1, slownik2):
    '''funkcja redukuje dwa slowniki z licznikami liter do jednego slownika'''
    wynik = slownik1
    for litera in slownik2.keys():
        if litera in wynik:
            wynik[litera] += slownik2[litera]
        else:
            wynik[litera] = slownik2[litera]
    return wynik

#tekst w ktorym policzymy wystapienia liter
lorem_ipsum = '''Lorem ipsum dolor sit amet, 
consectetur adipiscing elit, sed do eiusmod 
tempor incididunt ut labore et dolore magna aliqua. 
Ut enim ad minim veniam, quis nostrud exercitation 
ullamco laboris nisi ut aliquip ex ea commodo consequat. 
Duis aute irure dolor in reprehenderit in 
voluptate velit esse cillum dolore eu 
fugiat nulla pariatur. Excepteur sint 
occaecat cupidatat non proident, sunt in 
culpa qui officia deserunt mollit anim id est laborum.'''

#obliczenia beda wykonane niezaleznie na zdaniach (rozdzielonych kropkami)
zdania = lorem_ipsum.split('.')

#wykonujemy zadanie policz_litery na poszczegolnych zdaniach (mapowanie)
pulaworkerow = multiprocessing.Pool()
policzone_zdaniami = pulaworkerow.map(policz_litery,zdania)

#redukujemy otrzymane wyniki czesciowe przy pomocy funkcji polacz_slowniki
zredukowane = reduce(polacz_slowniki,policzone_zdaniami)
print(zredukowane)

Wieloprocesowy ping i xlsx

1Wieloprocesowy ping i obsługa xlsx w Python

Korzystając z biblioteki openpyxl można w pythonie całkiem wygodnie korzystać z plików w formacie xlsx. Poniższy programik ilustruje zarówno czytanie, jak i pisanie do xlsx-a.
Program jest też prostym przykładem wieloprocesowego przetważania. Wykorzystuje bibliotekę standardową multiprocessing, żeby uruchomić w procesach potomnych systemowe polecenie ping (przy pomocy subprocess.Popen). Hosty wczytane z pliku xlsx po kolei są ,,pingane'', a wyniki ping-a zapisywane są do pliku. Jeżeli część pakietów nie powróciła z celu, to komórka arkusza jest kolorowana.
'''program wysyła ping-a asynchronicznie do hostow
podanych jako argument skryptu i zapisuje do xlsx-a
wynik dla kazdego hosta.

Hosty do ktorych nie wszystkie ping-i doszly
zostana pokolorowane na czerwono.
'''

import multiprocessing
import re
import subprocess
import sys
from openpyxl import load_workbook, Workbook
from openpyxl.styles import PatternFill

def pingHost(hostname):
        try:
                v=subprocess.check_output('ping.exe -n 1 {0}'.format(hostname),shell=True)
        except subprocess.CalledProcessError as e:
                v=e.output              
        #zbierz lancuch zawierajacy wyniki pinga i podziel go na linie
        ping_output = v.decode(errors='ignore').splitlines()
        straty = '\n'.join([x.strip() for x in ping_output if re.search('straty',x)])
        ip = '\n'.join([re.sub('[^0-9.]','',x) for x in ping_output if re.search('badania',x)])
        return {'ip':ip,'straty':straty,'hostname':hostname}

def write_hosts(filename,host_list):
        '''zapisz otrzymane wyniki do xlsx'''
        wb = load_workbook(filename)
        ws = wb.active
        ws.cell(row=1,column=2).value='ip'
        ws.cell(row=1,column=3).value='losses'
        for i,host in enumerate(host_list):
                if(host['straty']!='(0% straty),'):
                        ws.cell(row=2+i,column=1).fill=red_pattern_fill()
                ws.cell(row=2+i,column=2).value=host['ip']
                ws.cell(row=2+i,column=3).value=host['straty']
        wb.save(filename)

def read_hosts(filename):
        '''wczytaj hosty do pingania z xlsx-a
        z pominieciem 1 wiersza (naglowka)'''
        wb = load_workbook(filename)
        ws = wb.active
        rows_to_read = range(2,ws.max_row+1)
        return [ws.cell(row=row,column=1).value for row in rows_to_read]

def red_pattern_fill():
        '''pomocnicza funkcja definiujaca czerwone wypelnienie'''
        return PatternFill(start_color='FFFF0000', end_color='FFFF0000', fill_type='solid')

if __name__ == '__main__':
        hosty=read_hosts(sys.argv[1])
        pool=multiprocessing.Pool()
        v=pool.map(pingHost,hosty)
        write_hosts(sys.argv[1],v)