Les nœuds PyODPS 3 vous permettent de rédiger des tâches MaxCompute en Python et de les planifier pour une exécution périodique dans DataWorks.
Présentation
PyODPS est le SDK Python pour MaxCompute. Il permet de rédiger des tâches, d'interroger des tables et des vues, ainsi que de gérer les ressources MaxCompute en Python. Pour plus d'informations, consultez la rubrique Présentation de PyODPS. Dans DataWorks, les nœuds PyODPS permettent de planifier et d'exécuter des tâches Python conjointement avec d'autres jobs.
Remarques sur l'utilisation
-
Pour appeler un package tiers depuis un nœud PyODPS sur un groupe de ressources DataWorks, utilisez un groupe de ressources serverless avec une image personnalisée.
RemarqueCette méthode ne s'applique pas si votre code inclut une fonction définie par l'utilisateur (UDF) qui fait référence à un package tiers. Pour la procédure correcte, consultez la rubrique Exemple UDF : Utiliser un package tiers dans une UDF Python.
Pour mettre à niveau la version de PyODPS, utilisez une image personnalisée dans un groupe de ressources serverless afin d'exécuter la commande
/home/tops/bin/pip3 install pyodps==0.12.1. Remplacez0.12.1par la version cible de PyODPS. Pour un groupe de ressources exclusif pour la planification, utilisez l'Assistant O&M pour exécuter la même commande.Si votre job PyODPS doit accéder à un environnement réseau spécifique, tel qu'une source de données ou un service dans un VPC ou un réseau IDC, utilisez un groupe de ressources serverless. Pour savoir comment connecter le groupe de ressources serverless à l'environnement cible, consultez la rubrique Solutions de connectivité réseau.
Pour plus d'informations sur la syntaxe PyODPS, consultez la documentation PyODPS.
Il existe deux types de nœuds PyODPS : PyODPS 2 (Python 2) et PyODPS 3 (Python 3). Créez le type de nœud correspondant à votre version de Python.
-
Si l'exécution d'instructions SQL dans un nœud PyODPS ne génère pas une lignée de données correcte dans Data Map, définissez manuellement les paramètres de planification DataWorks pertinents dans le code du job. Pour savoir comment afficher la lignée de données, consultez la rubrique Afficher la lignée de données. Pour savoir comment définir les paramètres, consultez la rubrique Définir les indicateurs de paramètres d'exécution. L'exemple de code suivant récupère les paramètres requis lors de l'exécution.
import os ... # get DataWorks scheduler runtime parameters skynet_hints = {} for k, v in os.environ.items(): if k.startswith('SKYNET_'): skynet_hints[k] = v ... # setting hints while submitting a job o.execute_sql('INSERT OVERWRITE TABLE XXXX SELECT * FROM YYYY WHERE ***', hints=skynet_hints) ... La taille maximale du journal de sortie d'un nœud PyODPS est de 4 Mo. Évitez d'imprimer de grandes quantités de données dans le journal. Concentrez-vous sur l'affichage des informations d'alerte et de progression.
Limitations
Lorsque vous exécutez un nœud PyODPS sur un groupe de ressources exclusif pour la planification, nous recommandons de ne pas dépasser 50 Mo de données traitées localement. Cette limite est imposée par les spécifications du groupe de ressources exclusif. Si les données locales dépassent le seuil du système d'exploitation, une erreur OOM (Got Killed) peut se produire. Évitez d'écrire un code de traitement de données excessif dans le nœud PyODPS.
-
Lorsque vous exécutez un nœud PyODPS sur un groupe de ressources serverless, configurez un nombre approprié d'unités de calcul (CU) en fonction du volume de données.
RemarqueSur un groupe de ressources serverless, un seul job prend en charge un maximum de
64 CUs. Toutefois, nous recommandons de ne pas dépasser16 CUsafin d'éviter les pénuries de ressources qui pourraient affecter le démarrage du job. Une erreur Got killed indique que le processus a dépassé la limite de mémoire. Évitez les opérations sur les données locales. Les jobs SQL et DataFrame (à l'exception des opérations to_pandas) initiés via PyODPS ne sont pas soumis à cette limitation.
Vous pouvez utiliser les bibliothèques NumPy et pandas préinstallées dans le code qui n'implique pas de fonctions définies par l'utilisateur. Les autres packages tiers contenant du code binaire ne sont pas pris en charge.
Pour des raisons de compatibilité,
options.tunnel.use_instance_tunnelest défini surFalsepar défaut dans DataWorks. Pour activerinstance tunnelglobalement, définissez manuellement cette valeur surTrue.-
La définition du bytecode diffère entre les versions mineures de Python 3, telles que Python 3.8 et Python 3.7.
MaxCompute utilise actuellement Python 3.7. Si vous utilisez une syntaxe issue d'autres versions de Python 3, comme le bloc
finally blockde Python 3.8, une erreur se produit lors de l'exécution. Nous vous recommandons d'utiliser Python 3.7. PyODPS 3 prend en charge l'exécution sur un groupe de ressources serverless. Pour en acheter et en utiliser un, consultez la rubrique Utiliser un groupe de ressources serverless.
Vous ne pouvez pas configurer l'exécution simultanée de plusieurs jobs Python au sein d'un seul nœud PyODPS.
Pour imprimer des journaux dans un nœud PyODPS, utilisez
print. L'utilisation delogger.infon'est pas prise en charge.
Prérequis
Associer un moteur de calcul MaxCompute à votre espace de travail DataWorks.
Procédure
-
Développez votre code sur la page de l'éditeur du nœud PyODPS 3.
Exemples de code PyODPS 3
Après avoir créé un nœud PyODPS, vous pouvez modifier et exécuter votre code. Pour plus d'informations sur la syntaxe PyODPS, consultez la rubrique Opérations de base. Les exemples suivants couvrent cinq scénarios courants. Choisissez celui qui correspond à vos besoins.
Point d'entrée ODPS
Chaque nœud PyODPS dans DataWorks inclut une variable globale de point d'entrée ODPS,
odpsouo, que vous n'avez pas besoin de définir manuellement.print(odps.exist_table('PyODPS_iris'))Exécuter SQL
Vous pouvez exécuter des instructions SQL dans un nœud PyODPS. Pour plus d'informations, consultez la rubrique SQL.
-
Par défaut,
instance tunnelest désactivé dans DataWorks, doncinstance.open_readerutilise l'interface Result et renvoie un maximum de 10 000 enregistrements. Vous pouvez utiliserreader.countpour obtenir le nombre d'enregistrements. Pour parcourir toutes les données, désactivez lalimit. Les instructions suivantes activent globalementinstance tunnelet désactivent lalimit.options.tunnel.use_instance_tunnel = True options.tunnel.limit_instance_tunnel = False # Disable the limit to read all data. with instance.open_reader() as reader: # All data can be read through Instance Tunnel. -
Vous pouvez également activer
instance tunnelpour un seul appelopen_readeren ajoutanttunnel=Trueà l'appelopen_reader. Vous pouvez aussi ajouterlimit=Falsepour désactiver la restrictionlimitpour cet appel.# Use Instance Tunnel for this open_reader call and read all data. with instance.open_reader(tunnel=True, limit=False) as reader:
Définir les paramètres d'exécution
-
Définissez les paramètres d'exécution à l'aide du paramètre
hints, qui est undict. Pour plus d'informations sur les hints, consultez la rubrique Opérations SET.o.execute_sql('select * from PyODPS_iris', hints={'odps.sql.mapper.split.size': 16}) -
Si vous définissez des configurations globales à l'aide de
sql.settings, les paramètres d'exécution pertinents sont ajoutés à chaque exécution.from odps import options options.sql.settings = {'odps.sql.mapper.split.size': 16} o.execute_sql('select * from PyODPS_iris') # Hints are added based on the global settings.
Lire les résultats d'exécution
Appelez
open_readerdirectement sur une instance d'exécution SQL. Cela prend en charge deux scénarios :-
Le SQL renvoie des données structurées.
with o.execute_sql('select * from dual').open_reader() as reader: for record in reader: # Process each record. -
Si vous exécutez des instructions SQL telles que
desc, utilisez la propriétéreader.rawpour obtenir les résultats bruts de l'exécution SQL.with o.execute_sql('desc dual').open_reader() as reader: print(reader.raw)RemarqueSi vous utilisez des paramètres de planification personnalisés et exécutez directement un nœud PyODPS 3 depuis l'interface utilisateur, vous devez coder en dur la valeur temporelle, car le nœud ne peut pas substituer les variables lors de l'exécution.
DataFrame
Vous pouvez également traiter les données à l'aide de DataFrame (Obsolète).
-
Exécution
Dans l'environnement DataWorks, les opérations DataFrame nécessitent un appel explicite à une méthode d'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(): # Call an immediate execution method to process each record.Si vous devez déclencher une exécution immédiate lors de l'impression, activez
options.interactive.from odps import options from odps.df import DataFrame options.interactive = True # Enable the option at the beginning. iris = DataFrame(o.get_table('pyodps_iris')) print(iris.sepal_width.sum()) # This triggers immediate execution. -
Imprimer des informations détaillées
options.verboseest activé par défaut dans DataWorks, ce qui entraîne l'impression d'informations détaillées telles que l'URL Logview lors de l'exécution.
Développement de code PyODPS 3
L'exemple suivant montre comment utiliser un nœud PyODPS :
Préparez un jeu de données en créant la table exemple pyodps_iris. Pour plus d'informations, consultez la rubrique Traiter les données DataFrame.
Créez un DataFrame. Pour plus d'informations, consultez la rubrique Créer un DataFrame à partir d'une table MaxCompute.
-
Saisissez le code suivant dans le nœud PyODPS et exécutez-le.
from odps.df import DataFrame # Create a DataFrame from an ODPS table. iris = DataFrame(o.get_table('pyodps_iris')) print(iris.sepallength.head(5))
Exécuter le job PyODPS****
-
Dans le volet Run Configuration, sous la section Compute Resource, configurez la Compute Resource, le computing quota et le DataWorks Resource Group.
RemarquePour accéder à une source de données via un réseau public ou un VPC, vous devez utiliser un groupe de ressources de planification ayant réussi un test de connectivité avec la source de données. Pour plus d'informations, consultez la rubrique Solutions de connectivité réseau.
Vous pouvez configurer les informations d'Image en fonction des exigences du job.
Dans la boîte de dialogue des paramètres de la barre d'outils, sélectionnez la source de données MaxCompute créée et cliquez sur Run pour exécuter le job PyODPS.
-
-
Pour exécuter le nœud périodiquement, configurez ses propriétés de planification en fonction de vos besoins métier. Pour plus d'informations, consultez la rubrique Configurer la planification des nœuds.
Contrairement aux nœuds SQL dans DataWorks, les nœuds PyODPS ne remplacent pas les chaînes telles que ${param_name} dans le code. Au lieu de cela, avant l'exécution du code, un dictionnaire nommé
argsest ajouté aux variables globales, à partir duquel vous pouvez récupérer les paramètres de planification. Par exemple, si vous définissezds=${yyyymmdd}dans la section Parameter, vous pouvez récupérer les informations de paramètre dans le code comme suit.print('ds=' + args['ds']) ds=20240930RemarqueSi vous devez obtenir la partition nommée
ds, vous pouvez utiliser la méthode suivante.o.get_table('table_name').get_partition('ds=' + args['ds']) Une fois le nœud configuré, vous devez le déployer. Pour plus d'informations, consultez la rubrique Déploiement des nœuds et des workflows.
Une fois le job déployé, vous pouvez afficher son statut dans Operation Center. Pour plus d'informations, consultez la rubrique Prise en main d'Operation Center.
Exécuter un nœud à l'aide d'un rôle associé
Vous pouvez associer un rôle RAM pour exécuter un nœud, ce qui vous permet d'exécuter des tâches de nœud avec un rôle RAM spécifique pour un contrôle granulaire des autorisations et une gestion sécurisée.
Étapes suivantes
FAQ sur PyODPS : Problèmes courants lors de l'exécution de PyODPS et comment les résoudre.