PyODPS DataFrame vous permet d'étendre les calculs intégrés grâce à des fonctions définies par l'utilisateur (UDF) et à des packages Python tiers. Cette rubrique aborde le mappage élément par élément avec map, les transformations au niveau des lignes et les agrégations personnalisées avec apply, les références aux ressources dans les UDF, ainsi que la procédure de chargement et de configuration des packages tiers.
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
Un objet DataFrame créé à partir d'une table MaxCompute ou d'un DataFrame pandas
La prise en charge des UDF Python activée dans votre projet MaxCompute (requise pour
mapetapplyavec des fonctions Python)
Les services publics Alibaba Cloud ne prennent pas en charge les UDF Python. Si votre projet ne prend pas en charge les UDF Python, la méthode map et les fonctions intégrées qui en dépendent ne sont pas disponibles.
Limites connues
| Limite | Détails |
|---|---|
| Types non pris en charge | Les méthodes map et apply n'acceptent ni les types LIST ni les types DICT en entrée ou en sortie. |
| Bibliothèque binaire préinstallée | NumPy est la seule bibliothèque tierce préinstallée contenant du code C. Toutes les autres bibliothèques binaires nécessitent un chargement explicite. |
| Compatibilité Python 2/3 | En raison des différences de bytecode entre les versions de Python, le code utilisant une syntaxe spécifique à Python 3 (telle que yield from) peut échouer sur un worker MaxCompute exécutant Python 2.7. Vérifiez que votre code s'exécute correctement avant de le déployer en production en utilisant l'API MapReduce dans Python 3. |
| Plateforme de compilation pour les packages binaires | Les fichiers wheel compilés sur macOS ou Windows ne peuvent pas être utilisés dans MaxCompute. Compilez les packages binaires dans un shell Linux. |
Appliquer des UDF à une colonne
Utilisez la méthode map sur un objet Sequence pour appeler une UDF sur chaque élément.
>>> 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
Si le type de la Sequence change après l'appel à map, spécifiez explicitement le nouveau type :
>>> 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
Éviter les bugs de capture de variables de fermeture
Lorsqu'une UDF contient une fermeture, les modifications externes apportées à la variable capturée affectent le comportement de la fonction. Le code suivant produit un résultat inattendu : chaque SequenceExpr dans dfs se retrouve sous la forme df.sepal_length + 9 :
>>> dfs = []
>>> for i in range(10):
>>> dfs.append(df.sepal_length.map(lambda x: x + i))
Corrigez ce problème en renvoyant la lambda depuis une fonction externe, ou en utilisant 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)))
Utiliser des UDF existantes
Transmettez un nom de fonction (chaîne) ou un objet Function à map pour appeler une UDF existante. Pour plus de détails, consultez la rubrique Functions.
Suivre l'exécution avec des compteurs
Utilisez get_execution_context pour accéder aux compteurs depuis une UDF. Les valeurs des compteurs apparaissent dans le JSONSummary de 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)
Appliquer des UDF à une ligne
Utilisez apply avec axis=1 pour appeler une UDF sur chaque ligne. L'UDF reçoit une ligne à la fois ; récupérez les valeurs des champs par nom d'attribut ou par index.
Renvoyer une seule valeur par ligne
Définissez reduce=True pour renvoyer une Sequence. Spécifiez le type de sortie avec le paramètre types (la valeur par défaut est 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
Renvoyer plusieurs lignes à l'aide de yield
Définissez reduce=False et utilisez yield pour émettre plusieurs lignes par ligne d'entrée. Spécifiez les noms et les types des champs de sortie avec names et 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
Annotez le schéma de sortie directement sur la fonction pour éviter de le répéter aux sites d'appel :
>>> 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
Équivalent : map_reduce en mode map uniquement
map_reduce en mode map uniquement est équivalent à apply avec axis=1 :
>>> iris.map_reduce(mapper=handle).count()
300
Utiliser une UDTF existante
Pour appeler une fonction table définie par l'utilisateur (UDTF) existante dans MaxCompute, transmettez son nom sous forme de chaîne :
>>> iris['name', 'sepallength'].apply('your_func', axis=1, names=['name2', 'sepallength2'], types=['string', 'float'])
Combiner la sortie de ligne avec une vue latérale
Lorsque reduce=False, combinez la sortie de l'UDF avec les colonnes d'origine à l'aide d'une vue latérale — utile pour les agrégations :
>>> 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)]
Appliquer des agrégations personnalisées à une colonne
Utilisez apply avec axis=0 (ou sans argument axis) pour transmettre une classe d'agrégation personnalisée sur tous les objets Sequence. La classe doit implémenter buffer, __call__, merge et 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
Lire des ressources MaxCompute dans les UDF
Les UDF peuvent lire des ressources MaxCompute — ressources de fichier et ressources de table — ou référencer un objet Collection en tant que ressource. Encapsulez l'UDF dans une fermeture ou une classe appelable afin que les ressources soient chargées une seule fois lors de l'initialisation plutôt qu'à chaque ligne.
Le chargement des ressources à l'intérieur de la fermeture (plutôt qu'à chaque appel de fonction) évite les frais généraux d'initialisation répétés — par exemple, lors du chargement de tables de recherche ou d'artefacts de modèle.
UDF au niveau de la ligne avec des ressources de fichier et de collection
>>> 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
Lors de la lecture de tables partitionnées, les champs de partition ne sont pas inclus.
UDF au niveau de la ligne avec un DataFrame local comme ressource
Les variables locales peuvent être référencées comme des ressources dans MaxCompute au moment de l'exécution. Dans l'exemple suivant, stop_words est un DataFrame local que l'exécuteur transmet à l'UDF en tant que ressource :
>>> 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
Pour les opérations sur les lignes (axis=1), utilisez une fermeture de fonction ou une classe appelable pour charger les ressources. Pour les agrégations de colonnes, utilisez plutôt la méthode__init__.
Charger des bibliothèques Python tierces
MaxCompute prend en charge le chargement de packages Python aux formats .whl, .egg, .zip et .tar.gz. Toutes les dépendances doivent être spécifiées explicitement ; l'omission d'une dépendance entraîne des erreurs d'importation au moment de l'exécution.
Choisissez votre méthode de chargement en fonction du type de package :
| Type de package | Méthode de chargement | Notes |
|---|---|---|
| Préinstallé | Aucun nécessaire | NumPy uniquement |
| Python pur (aucun code compilé, aucune opération sur les fichiers) | Charger en tant que ressource de fichier .whl |
Fonctionne pour des packages tels que python-dateutil, pytz, six. Les versions ultérieures de MaxCompute prennent également en charge les packages avec des opérations sur les fichiers. |
| Binaire (extensions C compilées) | Charger en tant que ressource d'archive .zip, activer l'isolation |
Nécessite le tag de plateforme cp27-cp27m-manylinux1_x86_64 ; compiler sur Linux |
Packages Python purs
Par défaut, PyODPS prend en charge les bibliothèques tierces contenant du code Python pur mais aucune opération sur les fichiers. L'exemple suivant charge python-dateutil et sa dépendance six.
Étape 1 : Téléchargez le package et ses dépendances. Les packages doivent être compilés pour Linux.
$ pip download python-dateutil -d /to/path/
Cela télécharge six-1.10.0-py2.py3-none-any.whl et python_dateutil-2.5.3-py2.py3-none-any.whl.
Étape 2 : Chargez les deux fichiers en tant que ressources à l'aide de 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'))
Étape 3 : Utilisez les bibliothèques. Spécifiez-les globalement via options.df.libraries, ou par exécution via le paramètre 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
Packages binaires (contenant du code compilé)
Les packages incluant des extensions C compilées (telles que SciPy ou pandas) nécessitent des étapes supplémentaires :
Le fichier
.whldoit utiliser le tag de plateformecp27-cp27m-manylinux1_x86_64.Chargez le fichier en tant que ressource d'archive, avec l'extension
.whlrenommée en.zip.Définissez
odps.isolation.session.enablesurTrue, ou activezisolationdans les paramètres de votre projet.
# 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'])
Pour compiler un package binaire à partir des sources, exécutez la commande suivante dans un shell Linux. Les fichiers wheel compilés sur macOS ou Windows ne sont pas compatibles avec MaxCompute.
python setup.py bdist_wheel
Chargement via la console MaxCompute
Comme alternative à l'API PyODPS, chargez les packages en utilisant add archive dans la console MaxCompute.
Étape 1 : Identifiez le fichier de package correct pour chaque dépendance.
La plupart des packages fournissent des fichiers .whl pour plusieurs plateformes. Pour les packages binaires, trouvez le fichier contenant cp27-cp27m-manylinux1_x86_64 dans son nom. Pour les packages Python purs, n'importe quel wheel py2.py3-none-any fonctionne.
Étape 2 : Vérifiez toutes les dépendances requises. Le tableau suivant répertorie les dépendances pour les packages courants.
| Package | Dépendances |
|---|---|
| pandas | NumPy, python-dateutil, pytz, six |
| SciPy | NumPy |
| scikit-learn | NumPy, SciPy |
NumPy est déjà préinstallé. Chargez uniquement python-dateutil, pytz, pandas, SciPy, scikit-learn et six.
Étape 3 : Téléchargez les fichiers de package. Le tableau suivant répertorie les fichiers spécifiques à télécharger pour chaque package.
| Package | Fichier à télécharger | Nom de la ressource chargée |
|---|---|---|
| python-dateutil | python-dateutil-2.6.0.zip | python-dateutil.zip |
| pytz | pytz-2017.2.zip | pytz.zip |
| six | six-1.11.0.tar.gz | six.tar.gz |
| pandas | pandas-0.20.2-cp27-cp27m-manylinux1_x86_64.zip | pandas.zip |
| SciPy | scipy-0.19.0-cp27-cp27m-manylinux1_x86_64.zip | scipy.zip |
| scikit-learn | scikit_learn-0.18.1-cp27-cp27m-manylinux1_x86_64.zip | sklearn.zip |
Étape 4 : Chargez chaque fichier. Pour les packages binaires (pandas, SciPy, scikit-learn), renommez l'extension .whl en .zip avant le chargement.
add archive python-dateutil.zip;
add archive pandas.zip;
Spécifier les bibliothèques pour l'exécution
Utilisez options.df.libraries pour définir les bibliothèques globalement pour la session, ou transmettez le paramètre libraries directement à une méthode d'exécution pour limiter la portée à un seul appel.
# 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'])