Les nœuds PyODPS 2 vous permettent d'écrire du code Python directement dans DataWorks pour traiter les données MaxCompute à l'aide du SDK Python PyODPS.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Créé un nœud PyODPS 2. Consultez la rubrique Créer et gérer des nœuds MaxCompute
Fonctionnement
PyODPS est le SDK Python pour MaxCompute. Dans DataWorks, chaque nœud PyODPS inclut une variable globale injectée automatiquement, odps (également aliasée sous le nom o), qui sert de point d'entrée vers MaxCompute. Utilisez ce point d'entrée pour écrire du code Python permettant d'interroger des tables, exécuter des instructions SQL, gérer des ressources et traiter des données. DataWorks injecte également un dictionnaire global args afin que votre code puisse lire les paramètres de planification lors de l'exécution.
Les nœuds PyODPS 2 et PyODPS 3 ne diffèrent que par la version de Python utilisée : PyODPS 2 s'exécute avec Python 2.7, tandis que PyODPS 3 utilise Python 3.
Limites
|**Contrainte**
|
**Détails**
| | --- | --- | |
**Version Python**
|
2.7
| |
**Traitement local des données (groupe de ressources exclusif)**
|
Maintenez la taille inférieure à 50 Mo. Le dépassement de cette limite peut provoquer une erreur de mémoire insuffisante (OOM) et entraîner l'arrêt du processus avec le message `Got killed`
| |
**Unités de calcul (groupe de ressources serverless)**
|
Jusqu'à 64 UC par tâche ; restez dans la limite de 16 UC pour éviter les pénuries de ressources au démarrage
| |
**Tâches Python simultanées**
|
Une seule tâche à la fois par nœud
| |
**Taille des journaux de sortie**
|
Jusqu'à 4 Mo
| |
**Packages tiers contenant du code binaire**
|
Non pris en charge
| |
**Bibliothèques préinstallées**
|
NumPy et pandas (utilisables en dehors des fonctions définies par l'utilisateur (UDF))
| |
**InstanceTunnel**
|
Désactivé par défaut ; définissez `options.tunnel.use_instance_tunnel = True` pour l'activer globalement
|
La limite de mémoire de 50 Mo s'applique uniquement aux opérations de données locales. Les tâches SQL et DataFrame (à l'exclusion de to_pandas ) initiées par PyODPS ne sont pas soumises à cette limite.
Remarques sur l'utilisation
Packages tiers
Pour utiliser un package tiers dans un nœud PyODPS, utilisez un groupe de ressources serverless et créez une image personnalisée qui inclut le package.
Si votre code UDF nécessite un package tiers, l'approche par image personnalisée ne s'applique pas. Consultez plutôt l'exemple : Référencer des packages tiers dans les UDF Python.
Accès réseau
Pour accéder à une source de données située dans un cloud privé virtuel (VPC) ou dans un centre de données sur site, exécutez le nœud sur un groupe de ressources serverless et configurez une connexion réseau entre le groupe de ressources et la source de données. Consultez les solutions de connectivité réseau.
Mise à niveau de PyODPS
Sur un groupe de ressources serverless : utilisez la fonctionnalité de gestion des images pour exécuter
/home/tops/bin/pip3 install pyodps==0.12.1sur un nœud PyODPS 3 (remplacez0.12.1par la version cible). Consultez la rubrique Gérer les images.Sur un groupe de ressources exclusif pour la planification : utilisez la fonctionnalité O&M Assistant pour exécuter la même commande sur un nœud PyODPS 3. Consultez la rubrique Utiliser la fonctionnalité O&M Assistant.
Journaux de sortie
Maintenez la sortie des journaux concise. Incluez les journaux d'alerte et les points de contrôle de progression plutôt que de déverser de grands ensembles de données dans le journal. La limite de 4 Mo s'applique à l'ensemble du journal de sortie d'une exécution de nœud.
Lignage des données
Si les instructions SQL dans un nœud PyODPS ne parviennent pas à générer des lignages de données dans Data Map, transmettez les paramètres d'exécution du planificateur DataWorks en tant qu'indicatifs SQL (hints). Le code suivant montre comment collecter ces paramètres et les transmettre lors de l'exécution du code SQL :
import os
# Collect DataWorks scheduler runtime parameters
skynet_hints = {}
for k, v in os.environ.items():
if k.startswith('SKYNET_'):
skynet_hints[k] = v
# Pass the parameters as hints when running SQL
o.execute_sql('INSERT OVERWRITE TABLE XXXX SELECT * FROM YYYY WHERE ***', hints=skynet_hints)
Pour plus d'informations, consultez les rubriques Afficher les lignages de données et Configurer le paramètre hints.
Écrire et exécuter du code
Pour la référence complète de la syntaxe PyODPS, consultez la rubrique Vue d'ensemble.
Utiliser le point d'entrée MaxCompute
Chaque nœud PyODPS injecte automatiquement odps (et son alias o) comme point d'entrée MaxCompute. Il n'est pas nécessaire d'initialiser manuellement un client.
print(odps.exist_table('PyODPS_iris'))
Exécuter des instructions SQL
Utilisez o.execute_sql() pour exécuter des instructions SQL sur MaxCompute.
Par défaut, InstanceTunnel est désactivé. Lors de la lecture des résultats avec instance.open_reader, l'interface Result limite la lecture à 10 000 enregistrements. Pour lire tous les enregistrements, activez InstanceTunnel globalement :
options.tunnel.use_instance_tunnel = True
options.tunnel.limit_instance_tunnel = False # Remove the record count limit
with instance.open_reader() as reader:
# Reads all records using InstanceTunnel
Pour activer InstanceTunnel pour une seule opération de lecture sans modifier le paramètre global :
with instance.open_reader(tunnel=True, limit=False) as reader:
# Reads all records for this operation only
Lire les résultats des requêtes SQL
Utilisez open_reader pour traiter les résultats des requêtes.
Pour les instructions SQL qui renvoient des données structurées :
with o.execute_sql('select * from dual').open_reader() as reader:
for record in reader: # Process each record
...
Pour les instructions DDL telles que DESC, utilisez reader.raw pour obtenir la sortie brute :
with o.execute_sql('desc dual').open_reader() as reader:
print(reader.raw)
Configurer les paramètres d'exécution
Transmettez les paramètres d'exécution à l'exécution SQL en utilisant le paramètre hints (un dictionnaire) :
o.execute_sql('select * from PyODPS_iris', hints={'odps.sql.mapper.split.size': 16})
Pour appliquer les paramètres globalement à toutes les exécutions SQL dans le nœud :
from odps import options
options.sql.settings = {'odps.sql.mapper.split.size': 16}
o.execute_sql('select * from PyODPS_iris') # Uses the global settings
Pour plus d'informations sur les indicatifs pris en charge, consultez la rubrique Opérations SET.
Utiliser DataFrame pour traiter les données
L'utilisation de DataFrame n'est pas recommandée. Envisagez d'utiliser SQL ou d'autres approches prises en charge.
Les opérations de l'API DataFrame sont différées : elles ne s'exécutent que lorsque vous appelez une méthode à exécution immédiate.
from odps.df import DataFrame
iris = DataFrame(o.get_table('pyodps_iris'))
for record in iris[iris.sepal_width < 3].execute(): # Triggers execution
...
Pour activer l'exécution implicite pour les méthodes d'affichage, définissez options.interactive sur True au début du nœud :
from odps import options
from odps.df import DataFrame
options.interactive = True # Enable at the top of the node
iris = DataFrame(o.get_table('pyodps_iris'))
print(iris.sepal_width.sum()) # Runs immediately and prints the result
Par défaut, options.verbose est réglé sur True dans DataWorks, de sorte que l'URL Logview et d'autres détails d'exécution apparaissent dans le journal lors d'une exécution.
Exemple : interroger une table MaxCompute avec DataFrame
Préparez un jeu de données et créez une table nommée
pyodps_iris. Consultez la rubrique Traitement des données DataFrame.Créez un objet DataFrame à partir de la table. Consultez la rubrique Créer un objet DataFrame à partir d'une table MaxCompute.
-
Saisissez le code suivant dans l'éditeur de code et exécutez le nœud :
from odps.df import DataFrame # Create a DataFrame from the MaxCompute table iris = DataFrame(o.get_table('pyodps_iris')) print(iris.sepallength.head(5))Sortie attendue :
sepallength 0 4.5 1 5.5 2 4.9 3 5.0 4 6.0
Configurer les paramètres de planification
Contrairement aux nœuds SQL, le code des nœuds PyODPS n'effectue pas de substitution de chaîne ${param_name} . Au lieu de cela, DataWorks injecte les paramètres de planification dans le nœud sous forme de dictionnaire global args avant l'exécution du nœud.
Étape 1 : Ajouter des paramètres dans les propriétés du nœud
Ouvrez l'onglet Properties du nœud (volet de navigation de droite dans l'onglet de configuration) et ajoutez des entrées dans la section Scheduling Parameter . Pour connaître les différences de format et de syntaxe entre les types de nœuds, consultez la rubrique Configurer les paramètres de planification pour différents types de nœuds.
Étape 2 : Lire les paramètres dans votre code
Par exemple, si vous définissez ds=${yyyymmdd} dans la section Scheduling Parameter , lisez la valeur dans votre code comme suit :
print('ds=' + args['ds'])
# Output: ds=20161116
Pour obtenir la partition correspondant à cette date :
o.get_table('table_name').get_partition('ds=' + args['ds'])
Les paramètres de planification personnalisés pour les nœuds PyODPS doivent être définis sur une valeur constante. Contrairement aux nœuds SQL, la valeur n'est pas automatiquement substituée lors de l'exécution.
Pour plus de scénarios de développement de tâches PyODPS, consultez :
Étapes suivantes
Vérifier que le nœud s'est exécuté avec succès : L'approche permettant de confirmer l'exécution réussie d'un script Shell s'applique également aux scripts Python.
Déployer le nœud PyODPS 3 : Dans les espaces de travail en mode standard, déployez le nœud dans l'environnement de production avant de le planifier.
Effectuer l'O&M sur le nœud PyODPS 3 : Après le déploiement dans Operation Center, gérez et surveillez le nœud depuis l'environnement de production.
FAQ PyODPS : Problèmes courants et conseils de dépannage.