Tous les produits
Search
Centre de documentation

MaxCompute:API MapReduce

Dernière mise à jour :Aug 10, 2026

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

Important

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

capacity

3000

Des valeurs plus élevées réduisent les faux positifs

Augmente

error_rate

0.01

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

columns

Champs à partir desquels extraire les paires clé-valeur

kv_delim

Délimiteur entre chaque clé et sa valeur

:

item_delim

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