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 dentro de UDFs e como enviar e configurar pacotes de terceiros.
Pré-requisitos
Antes de começar, verifique se você possui:
Um objeto DataFrame criado a partir de uma tabela do MaxCompute ou de um DataFrame do pandas
Suporte a UDF Python ativado em seu projeto do 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 que dependem dele estarã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 que contém 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 bytecode entre versões do Python, códigos que utilizam 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 pacotes binários em um shell Linux. |
Aplicar UDFs a uma coluna
Utilize 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
Caso o tipo da Sequence seja alterado após o uso de 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 erros 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 a seguir gera um resultado indesejado — cada SequenceExpr em dfs acaba sendo avaliado 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 utilizando 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)))
Utilizar UDFs existentes
Passe um nome de função (string) ou um objeto Function para map a fim de invocar uma UDF existente. Para mais detalhes, consulte Funções.
Monitorar execução com contadores
Use get_execution_context para acessar contadores de 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
Utilize 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 utilize 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 esquema 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 a 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 seu nome como uma string:
>>> iris['name', 'sepallength'].apply('your_func', axis=1, names=['name2', 'sepallength2'], types=['string', 'float'])
Combinar saída de linha com uma 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
Utilize 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 recursos de tabela — ou referenciar um objeto Collection como recurso. Envolva a UDF em uma closure ou classe chamável para que os recursos sejam carregados 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 um 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. Todas as dependências devem ser especificadas 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 compactado |
Requer a tag de plataforma |
Pacotes Python puro
Por padrão, o PyODPS suporta bibliotecas de terceiros que contêm código Python puro, mas nenhuma operação 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: Utilize 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 (contendo código compilado)
Pacotes que incluem 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 um recurso de arquivo compactado, renomeando a extensão
.whlpara.zip.Defina
odps.isolation.session.enablecomoTrueou ative oisolationnas configurações do seu 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 do PyODPS, envie pacotes usando o comando 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 múltiplas 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 a seguir 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 de enviar.
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'])