PyODPS propose des méthodes pour exécuter des instructions SQL MaxCompute et lire les résultats depuis Python.
|
Méthode |
Objectif |
|
|
Exécute le code SQL de manière synchrone (bloque jusqu'à la fin) |
|
|
Exécute le code SQL de manière asynchrone (renvoie immédiatement une instance) |
|
|
Lit les résultats de l'exécution SQL |
Toutes les instructions SQL ne peuvent pas être exécutées viaexecute_sql()etrun_sql(). Ces méthodes prennent en charge les instructions DDL (Data Definition Language) et DML (Data Manipulation Language). Utilisezrun_security_querypour les instructions GRANT ou REVOKE, etrun_xflowouexecute_xflowpour les appels API XFlow.
Exécuter des instructions SQL
Paramètres
|
Paramètre |
Type |
Description |
|
|
string |
L'instruction SQL à exécuter |
|
|
dict |
Paramètres d'exécution |
Valeurs de retour
Les méthodes execute_sql() et run_sql() renvoient toutes deux les informations relatives aux Instances de tâche.
Exemples
Exécution synchrone par rapport à asynchrone
# Synchronous: blocks until the statement finishes
o.execute_sql('select * from table_name')
# Asynchronous: returns immediately
instance = o.run_sql('select * from table_name')
print(instance.get_logview_address()) # Get the LogView URL
instance.wait_for_success() # Block until the statement finishes
Transmettre des hints d'exécution
o.execute_sql('select * from pyodps_iris', hints={'odps.stage.mapper.split.size': 16})
Paramètres globaux
Définissez options.sql.settings pour appliquer des paramètres d'exécution à chaque appel ultérieur de execute_sql(). Consultez les Paramètres Flag pour connaître les paramètres disponibles.
from odps import options
options.sql.settings = {'odps.stage.mapper.split.size': 16}
o.execute_sql('select * from pyodps_iris') # Hints apply automatically
Lire les résultats des requêtes
Appelez open_reader() sur une instance terminée pour lire les résultats. Le type de retour dépend de l'instruction SQL.
Données structurées (SELECT)
Les requêtes SELECT renvoient des enregistrements structurés. Utilisez une boucle for pour parcourir chaque enregistrement :
with o.execute_sql('select * from table_name').open_reader() as reader:
for record in reader:
print(record)
Données non structurées (DESC et autres commandes)
Des commandes telles que desc renvoient du texte brut. Accédez-y via reader.raw :
with o.execute_sql('desc table_name').open_reader() as reader:
print(reader.raw)
Interface InstanceTunnel par rapport à Result
Par défaut, open_reader() utilise l'interface Result, qui peut expirer ou limiter le nombre d'enregistrements renvoyés. Activez InstanceTunnel pour lire l'intégralité des données :
Option 1 : Paramètre global
from odps import options
options.tunnel.use_instance_tunnel = True
Option 2 : Paramètre par appel
with o.execute_sql('select * from table_name').open_reader(tunnel=True) as reader:
for record in reader:
print(record)
À partir de PyODPS V0.7.7.1, open_reader() prend en charge la lecture complète des données de cette manière.
Si votre version MaxCompute est plus ancienne ou si
InstanceTunnelrencontre une erreur, PyODPS génère une alerte et revient automatiquement à l'interfaceResult. Consultez le message d'alerte pour identifier la cause.Si votre version MaxCompute prend uniquement en charge l'interface
Resultet que vous avez besoin de tous les résultats, écrivez-les d'abord dans une autre table, puis lisez cette table avecopen_reader(). Cette approche est soumise au mécanisme de protection des données de votre projet.Pour plus d'informations sur
InstanceTunnel, consultez la section InstanceTunnel.
Mode de limitation de lecture
Par défaut, PyODPS ne limite pas les données lues depuis une instance. Toutefois, si le propriétaire du projet a configuré la protection des données, PyODPS active automatiquement le mode de limitation de lecture lorsque des restrictions sont détectées et que options.tunnel.limit_instance_tunnel n'est pas défini. Le mode de limitation de lecture plafonne généralement les lectures à 10 000 lignes.
Activer manuellement le mode de limitation de lecture (projets protégés) :
# Option A: per-call
reader = instance.open_reader(limit=True)
# Option B: global
options.tunnel.limit_instance_tunnel = True
Désactiver le mode de limitation de lecture pour lire toutes les données :
with o.execute_sql('select * from table_name').open_reader(tunnel=True, limit=False) as reader:
for record in reader:
print(record)
Dans des environnements tels que DataWorks,options.tunnel.limit_instance_tunnelpeut avoir la valeur par défautTrue. Pour lire toutes les données, transmettez à la foistunnel=Trueetlimit=Falseàopen_reader().
Projets protégés
Si votre projet est protégé et que tunnel=True, limit=False ne supprime pas la restriction, contactez le propriétaire du projet pour obtenir des autorisations de lecture. Consultez la section Protection des données du projet pour plus de détails.
Utiliser des alias de ressource
Lorsqu'une fonction définie par l'utilisateur (UDF) fait référence à des ressources qui changent dynamiquement, configurez un alias pour mapper l'ancien nom de ressource vers un nouveau. Cela évite de supprimer ou de recréer l'UDF.
Transmettez les alias via le paramètre aliases dans execute_sql() :
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, ], ]
# Write one row of data with value 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])
res2 = o.create_resource('test_alias_res2', 'file', file_obj='2')
# Map res1 alias to res2 without modifying the UDF or resource
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])
Configurer biz_id
Certains scénarios nécessitent un ID commercial (biz_id) pour l'exécution SQL. Si une erreur se produit en raison d'un biz_id manquant, définissez-le globalement :
from odps import options
options.biz_id = 'my_biz_id'
o.execute_sql('select * from pyodps_iris')