O PyODPS DataFrame permite estender a computação integrada com funções definidas pelo usuário (UDFs) e pacotes Python de terceiros. Este tópico aborda o mapeamento elemento a elemento com map, transformações no nível de linha e agregações personalizadas com apply, referências a recursos em UDFs e como enviar e configurar pacotes de terceiros.
Pré-requisitos
Antes de começar, verifique se você tem:
Um objeto DataFrame criado a partir de uma tabela do MaxCompute ou de um DataFrame pandas
Suporte a UDF Python ativado no projeto MaxCompute (necessário para usar
mapeapplycom funções Python)
Os serviços públicos da Alibaba Cloud não oferecem suporte a UDF Python. Se o seu projeto não tiver suporte a UDFs Python, o método map e as funções integradas dependentes ficarão indisponíveis.
Limitações conhecidas
|
Limitação |
Detalhes |
|
Tipos não suportados |
Os métodos |
|
Biblioteca binária pré-instalada |
A única biblioteca de terceiros pré-instalada com código C é a NumPy. Todas as outras bibliotecas binárias exigem envio explícito. |
|
Compatibilidade entre Python 2 e 3 |
Devido às diferenças de byte code entre versões do Python, códigos com sintaxe específica do Python 3 (como |
|
Plataforma de compilação para pacotes binários |
Arquivos Wheel compilados no macOS ou Windows não funcionam no MaxCompute. Compile os pacotes binários em um shell Linux. |
Aplicar UDFs a uma coluna
Use o método map em um objeto Sequence para chamar uma UDF em cada elemento.
>>> iris.sepallength.map(lambda x: x + 1).head(5)
sepallength
0 6.1
1 5.9
2 5.7
3 5.6
4 6.0
Se o tipo da Sequence mudar após o uso do map, especifique explicitamente o novo tipo:
>>> iris.sepallength.map(lambda x: 't' + str(x), 'string').head(5)
sepallength
0 t5.1
1 t4.9
2 t4.7
3 t4.6
4 t5.0
Evitar bugs de captura de variáveis de closure
Quando uma UDF contém uma closure, alterações externas na variável capturada afetam o comportamento da função. O código abaixo gera um resultado indesejado: cada SequenceExpr em dfs acaba interpretado como df.sepal_length + 9:
>>> dfs = []
>>> for i in range(10):
>>> dfs.append(df.sepal_length.map(lambda x: x + i))
Corrija esse problema retornando a lambda de uma função externa ou usando functools.partial:
# Option 1: use a factory function
>>> dfs = []
>>> def get_mapper(i):
>>> return lambda x: x + i
>>> for i in range(10):
>>> dfs.append(df.sepal_length.map(get_mapper(i)))
# Option 2: use functools.partial
>>> import functools
>>> dfs = []
>>> for i in range(10):
>>> dfs.append(df.sepal_length.map(functools.partial(lambda v, x: x + v, i)))
Usar UDFs existentes
Passe o nome da função (string) ou um objeto Function ao método map para invocar uma UDF existente. Para mais detalhes, consulte Functions.
Monitorar execução com contadores
Use get_execution_context para acessar contadores dentro de uma UDF. Os valores dos contadores aparecem no JSONSummary do LogView.
from odps.udf import get_execution_context
def h(x):
ctx = get_execution_context()
counters = ctx.get_counters()
counters.get_counter('df', 'add_one').increment(1)
return x + 1
df.field.map(h)
Aplicar UDFs a uma linha
Use apply com axis=1 para chamar uma UDF em cada linha. A UDF recebe uma linha por vez; recupere os valores dos campos pelo nome do atributo ou pelo índice.
Retornar um único valor por linha
Defina reduce=True para retornar uma Sequence. Especifique o tipo de saída com o parâmetro types (o padrão é STRING).
>>> iris.apply(lambda row: row.sepallength + row.sepalwidth, axis=1, reduce=True, types='float').rename('sepaladd').head(3)
sepaladd
0 8.6
1 7.9
2 7.9
Retornar múltiplas linhas usando yield
Configure reduce=False e use yield para emitir várias linhas para cada linha de entrada. Defina os nomes e tipos dos campos de saída com names e types.
>>> iris.count()
150
>>> def handle(row):
>>> yield row.sepallength - row.sepalwidth, row.sepallength + row.sepalwidth
>>> yield row.petallength - row.petalwidth, row.petallength + row.petalwidth
>>> iris.apply(handle, axis=1, names=['iris_add', 'iris_sub'], types=['float', 'float']).count()
300
Anote o schema de saída diretamente na função para evitar repeti-lo nos locais de chamada:
>>> from odps.df import output
>>> @output(['iris_add', 'iris_sub'], ['float', 'float'])
>>> def handle(row):
>>> yield row.sepallength - row.sepalwidth, row.sepallength + row.sepalwidth
>>> yield row.petallength - row.petalwidth, row.petallength + row.petalwidth
>>> iris.apply(handle, axis=1).count()
300
Equivalente: map_reduce apenas com map
O método map_reduce no modo apenas map equivale ao uso de apply com axis=1:
>>> iris.map_reduce(mapper=handle).count()
300
Usar uma UDTF existente
Para chamar uma função de tabela definida pelo usuário (UDTF) existente no MaxCompute, passe o nome dela como string:
>>> iris['name', 'sepallength'].apply('your_func', axis=1, names=['name2', 'sepallength2'], types=['string', 'float'])
Combinar saída de linha com lateral view
Quando reduce=False, combine a saída da UDF com as colunas originais usando uma lateral view — útil para agregações:
>>> from odps.df import output
>>> @output(['iris_add', 'iris_sub'], ['float', 'float'])
>>> def handle(row):
>>> yield row.sepallength - row.sepalwidth, row.sepallength + row.sepalwidth
>>> yield row.petallength - row.petalwidth, row.petallength + row.petalwidth
>>> iris[iris.category, iris.apply(handle, axis=1)]
Aplicar agregações personalizadas a uma coluna
Use apply com axis=0 (ou sem o argumento axis) para passar uma classe de agregação personalizada sobre todos os objetos Sequence. A classe deve implementar buffer, __call__, merge e getvalue.
class Agg(object):
def buffer(self):
return [0.0, 0]
def __call__(self, buffer, val):
buffer[0] += val
buffer[1] += 1
def merge(self, buffer, pbuffer):
buffer[0] += pbuffer[0]
buffer[1] += pbuffer[1]
def getvalue(self, buffer):
if buffer[1] == 0:
return 0.0
return buffer[0] / buffer[1]
>>> iris.exclude('name').apply(Agg)
sepallength_aggregation sepalwidth_aggregation petallength_aggregation petalwidth_aggregation
0 5.843333 3.054 3.758667 1.198667
Ler recursos do MaxCompute em UDFs
As UDFs podem ler recursos do MaxCompute — recursos de arquivo e de tabela — ou referenciar um objeto Collection como recurso. Envolva a UDF em uma closure ou classe chamável para carregar os recursos apenas uma vez na inicialização, em vez de a cada linha processada.
Carregar recursos dentro da closure (em vez de a cada chamada de função) evita sobrecarga de inicialização repetida — por exemplo, ao carregar tabelas de consulta ou artefatos de modelos.
UDF no nível de linha com recursos de arquivo e coleção
>>> file_resource = o.create_resource('pyodps_iris_file', 'file', file_obj='Iris-setosa')
>>> iris_names_collection = iris.distinct('name')[:2]
>>> iris_names_collection
sepallength
0 Iris-setosa
1 Iris-versicolor
>>> def myfunc(resources): # resources are passed in by calling order
>>> names = set()
>>> fileobj = resources[0] # file resources are represented by a file-like object
>>> for l in fileobj:
>>> names.add(l)
>>> collection = resources[1]
>>> for r in collection:
>>> names.add(r.name) # retrieve values by field name or offset
>>> def h(x):
>>> if x in names:
>>> return True
>>> else:
>>> return False
>>> return h
>>> df = iris.distinct('name')
>>> df = df[df.name,
>>> df.name.map(myfunc, resources=[file_resource, iris_names_collection], rtype='boolean').rename('isin')]
>>> df
name isin
0 Iris-setosa True
1 Iris-versicolor True
2 Iris-virginica False
Ao ler tabelas particionadas, os campos de partição não são incluídos.
UDF no nível de linha com DataFrame local como recurso
Variáveis locais podem ser referenciadas como recursos no MaxCompute durante a execução. No exemplo a seguir, stop_words é um DataFrame local que o executor passa para a UDF como recurso:
>>> words_df
sentence
0 Hello World
1 Hello Python
2 Life is short I use Python
>>> import pandas as pd
>>> stop_words = DataFrame(pd.DataFrame({'stops': ['is', 'a', 'I']}))
>>> @output(['sentence'], ['string'])
>>> def filter_stops(resources):
>>> stop_words = set([r[0] for r in resources[0]])
>>> def h(row):
>>> return ' '.join(w for w in row[0].split() if w not in stop_words),
>>> return h
>>> words_df.apply(filter_stops, axis=1, resources=[stop_words])
sentence
0 Hello World
1 Hello Python
2 Life short use Python
Para operações de linha (axis=1), use uma closure de função ou classe chamável para carregar recursos. Para agregações de coluna, utilize o método__init__.
Enviar bibliotecas Python de terceiros
O MaxCompute suporta o envio de pacotes Python nos formatos .whl, .egg, .zip e .tar.gz. Especifique todas as dependências explicitamente — omitir uma dependência causa erros de importação em tempo de execução.
Escolha o método de envio adequado ao tipo de pacote:
|
Tipo de pacote |
Método de envio |
Observações |
|
Pré-instalado |
Nenhum necessário |
Apenas NumPy |
|
Python puro (sem código compilado, sem operações de arquivo) |
Envie como recurso de arquivo |
Funciona para pacotes como python-dateutil, pytz, six. Versões mais recentes do MaxCompute também suportam pacotes com operações de arquivo. |
|
Binário (extensões C compiladas) |
Envie como recurso de arquivo |
Requer a tag de plataforma |
Pacotes Python puro
Por padrão, o PyODPS suporta bibliotecas de terceiros com código Python puro, mas sem operações de arquivo. O exemplo a seguir envia o python-dateutil e sua dependência six.
Etapa 1: Baixe o pacote e suas dependências. Os pacotes devem ser compilados para Linux.
$ pip download python-dateutil -d /to/path/
Isso baixa os arquivos six-1.10.0-py2.py3-none-any.whl e python_dateutil-2.5.3-py2.py3-none-any.whl.
Etapa 2: Envie ambos os arquivos como recursos usando create_resource.
# Make sure that file name extensions are correct.
>>> odps.create_resource('six.whl', 'file', file_obj=open('six-1.10.0-py2.py3-none-any.whl', 'rb'))
>>> odps.create_resource('python_dateutil.whl', 'file', file_obj=open('python_dateutil-2.5.3-py2.py3-none-any.whl', 'rb'))
Etapa 3: Use as bibliotecas. Especifique-as globalmente via options.df.libraries ou por execução através do parâmetro libraries.
# Global configuration (applies to all subsequent DataFrame operations in this session)
>>> from odps import options
>>> def get_year(t):
>>> from dateutil.parser import parse
>>> return parse(t).strftime('%Y')
>>> options.df.libraries = ['six.whl', 'python_dateutil.whl']
>>> df.datestr.map(get_year)
datestr
0 2016
1 2015
# Per-execution configuration (applies only to this call)
>>> def get_year(t):
>>> from dateutil.parser import parse
>>> return parse(t).strftime('%Y')
>>> df.datestr.map(get_year).execute(libraries=['six.whl', 'python_dateutil.whl'])
datestr
0 2016
1 2015
Pacotes binários (com código compilado)
Pacotes com extensões C compiladas (como SciPy ou pandas) exigem etapas adicionais:
O arquivo
.whldeve usar a tag de plataformacp27-cp27m-manylinux1_x86_64.Envie o arquivo como recurso de archive, renomeando a extensão
.whlpara.zip.Defina
odps.isolation.session.enablecomoTrueou ative oisolationnas configurações do projeto.
# Upload the binary package as an archive with the .zip extension.
>>> odps.create_resource('scipy.zip', 'archive', file_obj=open('scipy-0.19.0-cp27-cp27m-manylinux1_x86_64.whl', 'rb'))
# If isolation is already enabled in your project, the following option is optional.
>>> options.sql.settings = {'odps.isolation.session.enable': True}
>>> def psi(value):
>>> # Import the third-party library inside the function to avoid errors
>>> # caused by structural differences between operating systems.
>>> from scipy.special import psi
>>> return float(psi(value))
>>> df.float_col.map(psi).execute(libraries=['scipy.zip'])
Para compilar um pacote binário a partir do código-fonte, execute o comando abaixo em um shell Linux. Arquivos Wheel compilados no macOS ou Windows não são compatíveis com o MaxCompute.
python setup.py bdist_wheel
Envio via console do MaxCompute
Como alternativa à API PyODPS, envie pacotes usando add archive no console do MaxCompute.
Etapa 1: Identifique o arquivo de pacote correto para cada dependência.
A maioria dos pacotes fornece arquivos .whl para diversas plataformas. Para pacotes binários, localize o arquivo com cp27-cp27m-manylinux1_x86_64 no nome. Para pacotes Python puro, qualquer wheel py2.py3-none-any funciona.
Etapa 2: Verifique todas as dependências necessárias. A tabela a seguir lista as dependências dos pacotes mais comuns.
|
Pacote |
Dependências |
|
pandas |
NumPy, python-dateutil, pytz, six |
|
SciPy |
NumPy |
|
scikit-learn |
NumPy, SciPy |
O NumPy já vem pré-instalado. Envie apenas python-dateutil, pytz, pandas, SciPy, scikit-learn e six.
Etapa 3: Baixe os arquivos dos pacotes. A tabela abaixo lista os arquivos específicos a serem baixados para cada pacote.
|
Pacote |
Arquivo para baixar |
Nome do recurso de envio |
|
python-dateutil |
python-dateutil.zip |
|
|
pytz |
pytz.zip |
|
|
six |
six.tar.gz |
|
|
pandas |
pandas.zip |
|
|
SciPy |
scipy.zip |
|
|
scikit-learn |
sklearn.zip |
Etapa 4: Envie cada arquivo. Para pacotes binários (pandas, SciPy, scikit-learn), renomeie a extensão .whl para .zip antes do envio.
add archive python-dateutil.zip;
add archive pandas.zip;
Especificar bibliotecas para execução
Use options.df.libraries para definir bibliotecas globalmente na sessão ou passe o parâmetro libraries diretamente a um método de execução para restringi-lo a uma única chamada.
# Global: applies to all subsequent DataFrame operations in this session
>>> from odps import options
>>> options.df.libraries = ['six.whl', 'python_dateutil.whl']
# Local: applies only to this execution call
>>> df.apply(my_func, axis=1).to_pandas(libraries=['six.whl', 'python_dateutil.whl'])