L'API MapReduce de PyODPS DataFrame permet d'écrire une logique map et reduce personnalisée en Python et de l'exécuter à grande échelle sur MaxCompute. Une tâche map_reduce peut contenir uniquement des mappeurs, uniquement des réducteurs, ou les deux.
Exemple WordCount
L'exemple suivant compte les occurrences de mots dans une table contenant une seule colonne STRING.
#encoding=utf-8
from odps import ODPS
from odps import options
from odps.df import DataFrame
options.verbose = True
o = ODPS('your-access-id', 'your-secret-access-key',
project='DMP_UC_dev',
endpoint='http://service-corp.odps.aliyun-inc.com/api')
def mapper(row):
for word in row[0].split():
yield word.lower(), 1
def reducer(keys):
# Use a list instead of cnt=0. A plain integer would be treated as a local
# variable inside h(), so its value would not appear in the output.
cnt = [0]
def h(row, done): # done=True when all rows for this key have been processed
cnt[0] += row[1]
if done:
yield keys[0], cnt[0]
return h
word_count = DataFrame(o.get_table('zx_word_count'))
table = word_count.map_reduce(
mapper, reducer,
group=['word'],
mapper_output_names=['word', 'cnt'],
mapper_output_types=['string', 'int'],
reducer_output_names=['word', 'cnt'],
reducer_output_types=['string', 'int'],
)
Résultat attendu :
word cnt
0 are 1
1 day 1
2 doing? 1
...
Le paramètre group indique au réducteur le champ à utiliser pour regrouper les lignes entrantes. S'il est omis, tous les champs servent au regroupement. Le réducteur reçoit les keys agrégées et traite chaque ligne partageant ces clés. L'indicateur done prend la valeur True lorsque la dernière ligne d'une clé donnée a été traitée.
Il est possible d'écrire le réducteur sous forme de classe appelable plutôt que de fermeture :
class reducer(object):
def __init__(self, keys):
self.cnt = 0
def __call__(self, row, done): # done=True when all rows for this key are processed
self.cnt += row.cnt
if done:
yield row.word, self.cnt
Simplifier le schéma de sortie avec le décorateur @output
Le décorateur @output permet de déclarer directement les noms et types des champs de sortie sur la fonction, ce qui évite de transmettre mapper_output_names, mapper_output_types, reducer_output_names et reducer_output_types à map_reduce.
from odps.df import output
@output(['word', 'cnt'], ['string', 'int'])
def mapper(row):
for word in row[0].split():
yield word.lower(), 1
@output(['word', 'cnt'], ['string', 'int'])
def reducer(keys):
cnt = [0]
def h(row, done):
cnt[0] += row.cnt
if done:
yield keys.word, cnt[0]
return h
word_count = DataFrame(o.get_table('zx_word_count'))
table = word_count.map_reduce(mapper, reducer, group='word')
Pour trier les lignes au sein de chaque groupe pendant l'itération, transmettez le paramètre sort. Le paramètre ascending contrôle l'ordre de tri : une valeur booléenne unique applique le même ordre à tous les champs sort ; une liste permet de définir un ordre différent par champ (la longueur de la liste doit correspondre au nombre de champs sort).
Spécifier un combiner
Lors d'une tâche MapReduce, la sortie des mappeurs est transférée via le réseau vers les réducteurs — une étape appelée shuffle. Le shuffle implique des E/S disque, une sérialisation des données et un transfert réseau, ce qui en fait l'une des phases les plus coûteuses du pipeline.
Un combiner réduit le coût du shuffle en agrégeant localement la sortie des mappeurs sur chaque nœud avant son envoi aux réducteurs. Cette approche est particulièrement efficace pour les opérations commutatives et associatives, telles que le comptage et la sommation.
Contraintes avant d'écrire un combiner :
Un combiner ne peut pas référencer de ressources.
Les noms et types des champs de sortie doivent correspondre exactement à ceux du mappeur correspondant.
L'interface du combiner est identique à celle du réducteur. Transmettez la fonction combiner via le paramètre combiner :
words_df.map_reduce(mapper, reducer, combiner=reducer, group='word')
Référencer des ressources
Les mappeurs et les réducteurs peuvent chacun référencer leur propre ensemble de ressources. Les ressources sont transmises en tant que paramètres de fonction (par exemple, def mapper(resources):).
L'exemple suivant filtre les mots vides dans le mappeur et multiplie par 5 les comptes des mots figurant sur la liste blanche dans le réducteur :
white_list_file = o.create_resource('pyodps_white_list_words', 'file', file_obj='Python\nWorld')
@output(['word', 'cnt'], ['string', 'int'])
def mapper(resources):
stop_words = set(r[0].strip() for r in resources[0])
def h(row):
for word in row[0].split():
if word not in stop_words:
yield word, 1
return h
@output(['word', 'cnt'], ['string', 'int'])
def reducer(resources):
d = dict()
d['white_list'] = set(word.strip() for word in resources[0])
d['cnt'] = 0
def inner(keys):
d['cnt'] = 0
def h(row, done):
d['cnt'] += row.cnt
if done:
if row.word in d['white_list']:
d['cnt'] += 5
yield keys.word, d['cnt']
return h
return inner
words_df.map_reduce(mapper, reducer, group='word',
mapper_resources=[stop_words],
reducer_resources=[white_list_file])
Résultat attendu :
word cnt
0 hello 2
1 life 1
2 python 7
3 world 6
4 short 1
5 use 1
Utiliser des bibliothèques Python tierces
Les fonctionnalités de bytecode Python 3 telles que yield from provoquent des erreurs sur les workers MaxCompute exécutant Python 2.7. Testez votre code de bout en bout avant d'exécuter des tâches MapReduce basées sur Python 3 en production.
Spécifiez les bibliothèques globalement pour les appliquer à tous les appels map_reduce d'une session :
from odps import options
options.df.libraries = ['six.whl', 'python_dateutil.whl']
Ou transmettez-les à une seule exécution :
df.map_reduce(mapper=my_mapper, reducer=my_reducer, group='key').execute(
libraries=['six.whl', 'python_dateutil.whl']
)
Redistribuer les données
Lorsque les données sont distribuées de manière inégale entre les partitions du cluster, appelez reshuffle pour les rééquilibrer. Par défaut, les lignes sont attribuées aux partitions par hachage aléatoire :
df1 = df.reshuffle()
Pour distribuer selon une colonne spécifique et trier le résultat :
df1 = df.reshuffle('name', sort='id', ascending=False)
Filtre de Bloom
bloom_filter pré-filtre rapidement un jeu de données par rapport à un autre avant une jointure, réduisant ainsi le nombre de lignes à redistribuer et à comparer. Cette méthode fonctionne mieux lorsqu'un jeu de données est beaucoup plus grand que l'autre — par exemple, filtrer les données d'événements de navigation par rapport à un ensemble plus petit d'enregistrements de transactions avant de les joindre.
Le filtre est approximatif : il élimine les lignes absentes avec certitude de l'ensemble de référence, mais peut conserver un petit nombre de lignes qui n'y figurent pas réellement.
df1 = DataFrame(pd.DataFrame({'a': ['name1', 'name2', 'name3', 'name1'], 'b': [1, 2, 3, 4]}))
df2 = DataFrame(pd.DataFrame({'a': ['name1']}))
df1.bloom_filter('a', df2.a)
# The first argument can also be a column expression, e.g., df1.a + '1'
Résultat attendu :
a b
0 name1 1
1 name1 4
Dans cet exemple, name2 et name3 sont filtrés. Sur des jeux de données plus volumineux, le filtre peut ne pas éliminer toutes les lignes non correspondantes, mais le résultat de la jointure reste correct — le filtre n'affecte que les performances, pas la précision.
Configurez le filtre avec ces paramètres :
|
Paramètre |
Valeur par défaut |
Effet sur la précision |
Effet sur la mémoire |
|
|
|
Des valeurs plus élevées réduisent les faux positifs |
Augmente |
|
|
|
Des valeurs plus faibles réduisent les faux positifs |
Augmente |
Définissez les deux paramètres en fonction de la taille de votre jeu de données et de la mémoire disponible.
Pour plus d'informations sur l'exécution des opérations DataFrame, consultez Exécution DataFrame.
Tableau croisé dynamique
pivot_table résume les données en regroupant les lignes et en calculant des valeurs agrégées sur les colonnes.
Exemple de données :
>>> df
A B C D E
0 foo one small 1 3
1 foo one large 2 4
2 foo one large 2 5
3 foo two small 3 6
4 foo two small 3 4
5 bar one large 4 5
6 bar one small 5 3
7 bar two small 6 2
8 bar two large 7 1
Le paramètre rows est obligatoire. Il spécifie les champs de regroupement ; la fonction d'agrégation par défaut est mean :
>>> df['A', 'D', 'E'].pivot_table(rows='A')
A D_mean E_mean
0 bar 5.5 2.75
1 foo 2.2 4.40
Transmettez plusieurs champs à rows pour un regroupement plus fin :
>>> df.pivot_table(rows=['A', 'B', 'C'])
A B C D_mean E_mean
0 bar one large 4.0 5.0
1 bar one small 5.0 3.0
...
Utilisez values pour limiter les colonnes agrégées :
>>> df.pivot_table(rows=['A', 'B'], values='D')
A B D_mean
0 bar one 4.500000
1 bar two 6.500000
2 foo one 1.666667
3 foo two 3.000000
Utilisez aggfunc pour appliquer une ou plusieurs fonctions d'agrégation :
>>> df.pivot_table(rows=['A', 'B'], values=['D'], aggfunc=['mean', 'count', 'sum'])
A B D_mean D_count D_sum
0 bar one 4.500000 2 9
1 bar two 6.500000 2 13
2 foo one 1.666667 3 5
3 foo two 3.000000 2 6
Utilisez columns pour faire pivoter les valeurs d'une colonne en nouveaux en-têtes de colonne :
>>> df.pivot_table(rows=['A', 'B'], values='D', columns='C')
A B large_D_mean small_D_mean
0 bar one 4.0 5.0
1 bar two 7.0 6.0
2 foo one 2.0 1.0
3 foo two NaN 3.0
Utilisez fill_value pour remplacer les valeurs NaN par une valeur par défaut :
>>> df.pivot_table(rows=['A', 'B'], values='D', columns='C', fill_value=0)
A B large_D_mean small_D_mean
0 bar one 4 5
1 bar two 7 6
2 foo one 2 1
3 foo two 0 3
Conversion de chaînes clé-valeur
DataFrame peut analyser des chaînes clé-valeur en colonnes distinctes et convertir des données columnaires en chaînes clé-valeur. Pour plus d'informations sur la création d'objets DataFrame, consultez Créer un objet DataFrame.
Extraire les paires clé-valeur en colonnes
Utilisez extract_kv pour analyser une colonne contenant des chaînes clé-valeur délimitées :
>>> df
name kv
0 name1 k1=1,k2=3,k5=10
1 name1 k1=7.1,k7=8.2
2 name2 k2=1.2,k3=1.5
3 name2 k9=1.1,k2=1
>>> df.extract_kv(columns=['kv'], kv_delim='=', item_delim=',')
name kv_k1 kv_k2 kv_k3 kv_k5 kv_k7 kv_k9
0 name1 1.0 3.0 NaN 10.0 NaN NaN
1 name1 7.0 NaN NaN NaN 8.2 NaN
2 name2 NaN 1.2 1.5 NaN NaN NaN
3 name2 NaN 1.0 NaN NaN NaN 1.1
Paramètres :
|
Paramètre |
Description |
Valeur par défaut |
|
|
Champs à partir desquels extraire les paires clé-valeur |
— |
|
|
Délimiteur entre chaque clé et sa valeur |
|
|
|
Délimiteur entre les paires clé-valeur |
|
Les noms des colonnes de sortie suivent le modèle {original_field}_{key}, joints par un underscore. Les valeurs manquantes prennent par défaut la valeur NaN. Utilisez fill_value pour remplacer les valeurs manquantes par une valeur par défaut spécifique :
>>> df.extract_kv(columns=['kv'], kv_delim='=', fill_value=0)
name kv_k1 kv_k2 kv_k3 kv_k5 kv_k7 kv_k9
0 name1 1.0 3.0 0.0 10.0 0.0 0.0
1 name1 7.0 0.0 0.0 0.0 8.2 0.0
2 name2 0.0 1.2 1.5 0.0 0.0 0.0
3 name2 0.0 1.0 0.0 0.0 0.0 1.1
Convertir des colonnes en chaînes clé-valeur
Utilisez to_kv pour sérialiser plusieurs colonnes en une seule colonne de chaîne clé-valeur :
>>> df
name k1 k2 k3 k5 k7 k9
0 name1 1.0 3.0 NaN 10.0 NaN NaN
1 name1 7.0 NaN NaN NaN 8.2 NaN
2 name2 NaN 1.2 1.5 NaN NaN NaN
3 name2 NaN 1.0 NaN NaN NaN 1.1
>>> df.to_kv(columns=['k1', 'k2', 'k3', 'k5', 'k7', 'k9'], kv_delim='=')
name kv
0 name1 k1=1,k2=3,k5=10
1 name1 k1=7.1,k7=8.2
2 name2 k2=1.2,k3=1.5
3 name2 k9=1.1,k2=1