O PyODPS fornece métodos para executar instruções SQL do MaxCompute e ler resultados em Python.
|
Método |
Finalidade |
|
|
Executa SQL de forma síncrona (bloqueia até a conclusão) |
|
|
Executa SQL de forma assíncrona (retorna uma instância imediatamente) |
|
|
Lê os resultados da execução SQL |
Nem todas as instruções SQL são compatíveis comexecute_sql()erun_sql(). Esses métodos aceitam instruções de Linguagem de Definição de Dados (DDL) e Linguagem de Manipulação de Dados (DML). Userun_security_querypara instruções GRANT ou REVOKE erun_xflowouexecute_xflowpara chamadas à API XFlow.
Executar instruções SQL
Parâmetros
|
Parâmetro |
Tipo |
Descrição |
|
|
string |
Instrução SQL a executar |
|
|
dict |
Parâmetros de tempo de execução |
Valores de retorno
Tanto execute_sql() quanto run_sql() retornam informações de Instâncias de tarefa.
Exemplos
Execução síncrona versus assíncrona
# 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
Passar hints de tempo de execução
o.execute_sql('select * from pyodps_iris', hints={'odps.stage.mapper.split.size': 16})
Configurações globais
Defina options.sql.settings para aplicar parâmetros de tempo de execução a todas as chamadas subsequentes de execute_sql(). Consulte Parâmetros de flag para ver os parâmetros disponíveis.
from odps import options
options.sql.settings = {'odps.stage.mapper.split.size': 16}
o.execute_sql('select * from pyodps_iris') # Hints apply automatically
Ler resultados de consultas
Chame open_reader() em uma instância concluída para ler os resultados. O tipo de retorno varia conforme a instrução SQL.
Dados estruturados (SELECT)
Consultas SELECT retornam registros estruturados. Use um loop for para iterar sobre cada registro:
with o.execute_sql('select * from table_name').open_reader() as reader:
for record in reader:
print(record)
Dados não estruturados (DESC e outros comandos)
Comandos como desc retornam texto bruto. Acesse-o por meio de reader.raw:
with o.execute_sql('desc table_name').open_reader() as reader:
print(reader.raw)
Interface InstanceTunnel versus Result
Por padrão, open_reader() usa a interface Result, que pode expirar ou limitar o número de registros retornados. Ative o InstanceTunnel para leituras completas de dados:
Opção 1: Configuração global
from odps import options
options.tunnel.use_instance_tunnel = True
Opção 2: Parâmetro por chamada
with o.execute_sql('select * from table_name').open_reader(tunnel=True) as reader:
for record in reader:
print(record)
No PyODPS V0.7.7.1 e versões posteriores, open_reader() permite a leitura completa de dados dessa maneira.
Se a versão do MaxCompute for antiga ou se o
InstanceTunnelapresentar erro, o PyODPS gera um alerta e reverte automaticamente para a interfaceResult. Verifique a mensagem de alerta para identificar a causa.Caso a versão do MaxCompute suporte apenas a interface
Resulte todos os resultados sejam necessários, grave-os primeiro em outra tabela e leia dessa tabela comopen_reader(). Essa abordagem está sujeita ao mecanismo de proteção de dados do projeto.Para mais informações sobre o
InstanceTunnel, consulte InstanceTunnel.
Modo de limite de leitura
Por padrão, o PyODPS não limita a leitura de dados de uma instância. Contudo, se o proprietário do projeto configurar a proteção de dados, o PyODPS ativa automaticamente o modo de limite de leitura ao detectar restrições, desde que options.tunnel.limit_instance_tunnel não esteja definido. Na maioria dos casos, esse modo restringe a leitura a 10.000 linhas.
Ativar manualmente o modo de limite de leitura (projetos protegidos):
# Option A: per-call
reader = instance.open_reader(limit=True)
# Option B: global
options.tunnel.limit_instance_tunnel = True
Desativar o modo de limite de leitura para ler todos os dados:
with o.execute_sql('select * from table_name').open_reader(tunnel=True, limit=False) as reader:
for record in reader:
print(record)
Em ambientes como o DataWorks,options.tunnel.limit_instance_tunnelpode terTruecomo valor padrão. Para ler todos os dados, passetunnel=Trueelimit=Falseparaopen_reader().
Projetos protegidos
Se o projeto for protegido e tunnel=True, limit=False não remover a restrição, entre em contato com o proprietário do projeto para solicitar permissões de leitura. Consulte Proteção de dados do projeto para obter detalhes.
Usar aliases de recursos
Quando uma função definida pelo usuário (UDF) referencia recursos dinâmicos, configure um alias para mapear o nome antigo do recurso para o novo. Isso evita excluir ou recriar a UDF.
Passe aliases pelo parâmetro aliases em 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])
Configurar biz_id
Alguns cenários exigem um ID de negócio (biz_id) para executar SQL. Se ocorrer um erro por falta de biz_id, defina-o globalmente:
from odps import options
options.biz_id = 'my_biz_id'
o.execute_sql('select * from pyodps_iris')