PyODPS prend en charge les requêtes MaxCompute SQL en modes synchrone et asynchrone, et propose des méthodes pour lire les résultats sous forme d'enregistrements, de sortie brute ou de DataFrames pandas. Ce guide explique comment soumettre des instructions SQL depuis du code Python et récupérer efficacement les résultats.
Notes d'utilisation
Utilisez
execute_sql('statement')pour une exécution synchrone etrun_sql('statement')pour une exécution asynchrone. Les deux méthodes renvoient une instance de tâche en cours d'exécution.MaxCompute ne permet pas de lire les résultats d'instance au format Arrow.
-
Certaines instructions SQL fonctionnant dans la console MaxCompute ne sont pas exécutables dans PyODPS. Pour les instructions autres que DDL (Data Definition Language) et DML (Data Manipulation Language), utilisez la méthode appropriée :
run_security_query— exécutez les instructions GRANT et REVOKErun_xflowouexecute_xflow— exécutez les instructions Machine Learning Platform for AI (PAI)
La facturation des jobs SQL dépend du nombre de jobs soumis. Consultez la rubrique Présentation de la facturation pour plus de détails.
Exécuter des instructions SQL
Les méthodes execute_sql() et run_sql() soumettent toutes deux du code SQL à MaxCompute. La différence réside dans l'attente ou non du résultat par votre code.
import os
from odps import ODPS
# Store credentials in environment variables rather than hardcoding them
o = ODPS(
os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
project='<your-project>',
endpoint='<your-endpoint>',
)
# Synchronous: blocks until the SQL job completes
o.execute_sql('select * from table_name')
# Asynchronous: returns immediately; the job runs in the background
instance = o.run_sql('select * from table_name')
print(instance.get_logview_address()) # Get the Logview URL for monitoring
instance.wait_for_success() # Block until the job finishes
La méthode synchrone bloque le thread actuel jusqu'à la fin du job. La méthode asynchrone renvoie immédiatement une instance de tâche, ce qui permet de soumettre plusieurs jobs en parallèle et d'attendre leur achèvement de manière sélective.
Lors de l'exécution de requêtes asynchrones en parallèle, assurez-vous de comprendre l'ordre de dépendance entre les requêtes avant de les soumettre. Isolez les tâches asynchrones de longue durée dans des connexions distinctes si elles doivent survivre à la session actuelle.
Définir les paramètres d'exécution
Transmettez un paramètre hints sous forme de dict pour remplacer les paramètres de session pour une seule exécution :
o.execute_sql('select * from pyodps_iris', hints={'odps.sql.mapper.split.size': 16})
Pour appliquer les mêmes paramètres à chaque exécution de la session, configurez globalement options.sql.settings :
from odps import options
options.sql.settings = {'odps.sql.mapper.split.size': 16}
o.execute_sql('select * from pyodps_iris') # Picks up the global hints automatically
Lire les résultats des requêtes
Utilisez open_reader pour accéder aux résultats après l'achèvement d'une requête. Choisissez l'approche en fonction du type de retour de la requête et du volume de données requis.
|
Scénario |
Approche |
|
Données structurées (jeu de résultats petit à moyen) |
Parcourez les enregistrements avec |
|
Sortie texte brute (par exemple, commande |
Lisez |
|
Charger tous les résultats en mémoire |
|
|
Jeu de résultats volumineux avec téléchargement parallèle |
|
Données structurées — parcourez les enregistrements :
with o.execute_sql('select * from table_name').open_reader() as reader:
for record in reader:
print(record)
Sortie brute (par exemple, issue d'une commande DESC) — lisez reader.raw :
with o.execute_sql('desc table_name').open_reader() as reader:
print(reader.raw)
Choisir entre Instance Tunnel et l'interface Result
open_reader prend en charge deux backends :
|
Backend |
Quand l'utiliser |
Comment sélectionner |
|
Instance Tunnel |
Activation volontaire ; prend en charge les jeux de résultats volumineux |
Définissez |
|
Interface Result |
Valeur par défaut, ou solution de secours pour les anciennes versions de MaxCompute ou en cas d'erreurs Instance Tunnel |
|
PyODPS revient automatiquement à l'interface Result si Instance Tunnel échoue et enregistre une alerte expliquant la raison. Pour changer de backend manuellement :
# Use Instance Tunnel explicitly
with o.execute_sql('select * from dual').open_reader(tunnel=True) as reader:
for record in reader:
print(record)
# Use the Result interface explicitly
with o.execute_sql('select * from dual').open_reader(tunnel=False) as reader:
for record in reader:
print(record)
Limiter le nombre d'enregistrements téléchargés
Pour limiter le nombre d'enregistrements renvoyés, transmettez une option limit à open_reader, ou définissez options.tunnel.limit_instance_tunnel = True. Sans configuration explicite, MaxCompute applique la limite Tunnel au niveau du projet, généralement 10 000 enregistrements par téléchargement.
Lire les résultats dans un DataFrame pandas
Appelez to_pandas() sur le lecteur pour charger tous les résultats dans un DataFrame pandas :
with o.execute_sql('select * from dual').open_reader(tunnel=True) as reader:
pd_df = reader.to_pandas() # Returns a pandas DataFrame
Accélérer la lecture avec plusieurs processus
La lecture multiprocessus nécessite PyODPS 0.11.3 ou version ultérieure.
Transmettez n_process à to_pandas() pour distribuer le téléchargement sur plusieurs processus. Le SDK divise le jeu de résultats en lots et les télécharge en parallèle.
import multiprocessing
n_process = multiprocessing.cpu_count()
with o.execute_sql('select * from dual').open_reader(tunnel=True) as reader:
pd_df = reader.to_pandas(n_process=n_process)
Configurer les alias de ressources pour les UDF
Lorsque le code SQL exécute une fonction définie par l'utilisateur (UDF) qui référence un fichier de ressource et que le contenu de la ressource change, vous devez normalement supprimer et recréer l'UDF. Utilisez plutôt le paramètre aliases pour mapper l'ancien nom de ressource au nouveau : aucune modification de l'UDF n'est nécessaire.
L'exemple suivant configure une UDF qui lit un entier à partir d'un fichier de ressource et l'ajoute à la valeur d'entrée :
from odps.models import Schema
myfunc = '''\
from odps.udf import annotate
from odps.distcache import get_cache_file
@annotate('bigint->bigint')
class Example(object):
def __init__(self):
self.n = int(get_cache_file('test_alias_res1').read())
def evaluate(self, arg):
return arg + self.n
'''
res1 = o.create_resource('test_alias_res1', 'file', file_obj='1')
o.create_resource('test_alias.py', 'py', file_obj=myfunc)
o.create_function('test_alias_func',
class_type='test_alias.Example',
resources=['test_alias.py', 'test_alias_res1'])
table = o.create_table(
'test_table',
schema=Schema.from_lists(['size'], ['bigint']),
if_not_exists=True
)
data = [[1, ], ]
o.write_table(table, 0, [table.new_record(it) for it in data])
with o.execute_sql(
'select test_alias_func(size) from test_table').open_reader() as reader:
print(reader[0][0])
Pour exécuter la même requête avec une ressource différente sans modifier l'UDF, transmettez aliases :
res2 = o.create_resource('test_alias_res2', 'file', file_obj='2')
# Map the old resource name to the new one at query time
with o.execute_sql(
'select test_alias_func(size) from test_table',
aliases={'test_alias_res1': 'test_alias_res2'}).open_reader() as reader:
print(reader[0][0])
Exécuter du code SQL dans des environnements interactifs
PyODPS inclut des plug-ins SQL pour IPython et Jupyter qui prennent en charge les requêtes paramétrées. Consultez la documentation sur l'amélioration de l'expérience utilisateur pour obtenir les instructions de configuration.
Définir biz_id
Dans certains cas, vous devez transmettre biz_id lors de la soumission d'instructions SQL. Sinon, une erreur se produit lors de l'exécution. Définissez-le globalement afin qu'il s'applique à toutes les exécutions de la session :
from odps import options
options.biz_id = 'my_biz_id'
o.execute_sql('select * from pyodps_iris')