Todos os produtos
Search
Central de documentação

MaxCompute:API MapReduce

Última atualização: Jun 26, 2026

A API MapReduce do PyODPS DataFrame permite escrever lógicas personalizadas de map e reduce em Python e executá-las em escala no MaxCompute. Um job map_reduce pode conter apenas mappers, apenas reducers ou ambos.

Exemplo de WordCount

O exemplo a seguir conta as ocorrências de palavras em uma tabela que contém uma única coluna 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'],
)

Saída esperada:

     word  cnt
0     are    1
1     day    1
2  doing?    1
...

O parâmetro group define qual campo o reducer utiliza para agrupar as linhas recebidas. Caso seja omitido, todos os campos serão usados no agrupamento. O reducer recebe as keys agregadas e processa cada linha que compartilha essas chaves. O sinalizador done assume o valor True quando a última linha de uma determinada chave é processada.

É possível implementar o reducer como uma classe chamável (callable class) em vez de um closure:

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

Simplifique o schema de saída com o decorador @output

Utilize o decorador @output para declarar os nomes e tipos dos campos de saída diretamente na função. Isso elimina a necessidade de passar mapper_output_names, mapper_output_types, reducer_output_names e reducer_output_types para 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')

Para ordenar as linhas dentro de cada grupo durante a iteração, passe o parâmetro sort. O parâmetro ascending controla a direção da ordenação: um único valor booleano aplica a mesma ordem a todos os campos de sort; já uma lista permite definir uma direção diferente por campo (o tamanho da lista deve corresponder ao número de campos de sort).

Especifique um combiner

Durante um job MapReduce, a saída do mapper é transferida pela rede para os reducers — etapa conhecida como shuffle. Esse processo envolve E/S de disco, serialização de dados e transferência de rede, sendo uma das partes mais custosas do pipeline.

Um combiner reduz o custo do shuffle ao agregar a saída do mapper localmente em cada nó antes do envio aos reducers. Essa abordagem é mais eficaz para operações comutativas e associativas, como contagem e soma.

Restrições antes de escrever um combiner:

  • Um combiner não pode referenciar resources.

  • Os nomes e tipos dos campos de saída devem corresponder exatamente aos do mapper associado.

A interface do combiner é idêntica à do reducer. Passe a função do combiner por meio do parâmetro combiner:

words_df.map_reduce(mapper, reducer, combiner=reducer, group='word')

Referencie resources

Mappers e reducers podem referenciar seus próprios conjuntos de resources. Os resources são passados como parâmetros da função (por exemplo, def mapper(resources):).

O exemplo abaixo filtra stop words no mapper e incrementa em 5 a contagem de palavras da lista de permissões no reducer:

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])

Saída esperada:

    word  cnt
0  hello    2
1   life    1
2  python    7
3   world    6
4   short    1
5     use    1

Utilize bibliotecas Python de terceiros

Importante

Recursos de bytecode do Python 3, como yield from, causam erros em workers do MaxCompute que executam Python 2.7. Teste seu código de ponta a ponta antes de executar jobs MapReduce baseados em Python 3 em produção.

Especifique as bibliotecas globalmente para aplicá-las a todas as chamadas de map_reduce em uma sessão:

from odps import options
options.df.libraries = ['six.whl', 'python_dateutil.whl']

Ou passe-as para uma única execução:

df.map_reduce(mapper=my_mapper, reducer=my_reducer, group='key').execute(
    libraries=['six.whl', 'python_dateutil.whl']
)

Reembaralhe dados

Quando os dados estiverem distribuídos de forma desigual entre as partições do cluster, chame reshuffle para reequilibrá-los. Por padrão, as linhas são atribuídas às partições por hash aleatório:

df1 = df.reshuffle()

Para distribuir por uma coluna específica e ordenar o resultado:

df1 = df.reshuffle('name', sort='id', ascending=False)

Filtro Bloom

O bloom_filter pré-filtra rapidamente um conjunto de dados em relação a outro antes de um join, reduzindo a quantidade de linhas que precisam ser embaralhadas e comparadas. Ele funciona melhor quando um conjunto de dados é muito maior que o outro — por exemplo, filtrar dados de eventos de navegação contra um conjunto menor de registros de transações antes de uni-los.

Esse filtro é aproximado: ele elimina linhas definitivamente ausentes no conjunto de referência, mas pode reter um pequeno número de linhas que não estão realmente presentes.

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'

Saída esperada:

       a  b
0  name1  1
1  name1  4

Neste exemplo, name2 e name3 são filtrados. Em conjuntos de dados maiores, o filtro pode não eliminar todas as linhas sem correspondência, mas o resultado do join permanece correto — o filtro afeta apenas o desempenho, não a precisão.

Configure o filtro com estes parâmetros:

Parâmetro

Padrão

Efeito na precisão

Efeito na memória

capacity

3000

Valores maiores reduzem falsos positivos

Aumenta

error_rate

0.01

Valores menores reduzem falsos positivos

Aumenta

Defina ambos os parâmetros com base no tamanho do seu conjunto de dados e na memória disponível.

Para obter informações sobre a execução de operações de DataFrame, consulte Execução de DataFrame.

Tabela dinâmica

A função pivot_table resume dados agrupando linhas e calculando valores agregados nas colunas.

Dados de amostra:

>>> 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

O parâmetro rows é obrigatório. Ele especifica os campos para agrupamento; a função de agregação padrão é 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

Passe múltiplos campos para rows para obter um agrupamento mais granular:

>>> 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
...

Use values para restringir as colunas que serão agregadas:

>>> 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

Utilize aggfunc para aplicar uma ou mais funções de agregação:

>>> 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

Empregue columns para transformar os valores de uma coluna em novos cabeçalhos de coluna:

>>> 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

Aplique fill_value para substituir NaN por um valor padrão:

>>> 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

Conversão de strings chave-valor

O DataFrame consegue analisar strings chave-valor em colunas separadas e converter dados colunares de volta para strings chave-valor. Para obter informações sobre a criação de objetos DataFrame, consulte Crie um objeto DataFrame.

Extraia pares chave-valor para colunas

Use extract_kv para analisar uma coluna que contém strings chave-valor delimitadas:

>>> 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

Parâmetros:

Parâmetro

Descrição

Padrão

columns

Campos dos quais extrair pares chave-valor

kv_delim

Delimitador entre cada chave e seu valor

:

item_delim

Delimitador entre pares chave-valor

,

Os nomes das colunas de saída seguem o padrão {original_field}_{key}, unidos por um underscore. Valores ausentes assumem o padrão NaN. Utilize fill_value para substituir valores ausentes por um padrão específico:

>>> 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

Converta colunas para strings chave-valor

Utilize to_kv para serializar várias colunas em uma única coluna de string chave-valor:

>>> 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