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
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 |
|
|
|
Valores maiores reduzem falsos positivos |
Aumenta |
|
|
|
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 |
|
|
Campos dos quais extrair pares chave-valor |
— |
|
|
Delimitador entre cada chave e seu valor |
|
|
|
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