O PyODPS DataFrame usa execução adiada: as operações só rodam quando acionadas explicitamente. Este tópico explica como acionar a execução, recuperar e salvar resultados, configurar parâmetros de tempo de execução e executar operações de forma assíncrona ou em paralelo.
Pré-requisitos
Antes de começar, verifique se você tem:
Uma tabela de exemplo chamada
pyodps_iris. Consulte Processamento de dados do DataFrame para obter instruções de configuração.Um objeto DataFrame criado a partir de uma tabela do MaxCompute. Consulte a seção "Criar um objeto DataFrame a partir de uma tabela do MaxCompute" em Criar um objeto DataFrame.
Como funciona
O PyODPS DataFrame separa a definição de uma operação da sua execução. Ao encadear filtros, projeções ou agregações, nada é executado; você constrói um plano lógico. A execução começa apenas quando você chama uma ação, ou seja, um método que aciona explicitamente o plano e retorna resultados.
Os métodos a seguir são ações:
|
Método |
Descrição |
Retorno |
|
|
Executa a operação e retorna todos os resultados |
ResultFrame |
|
|
Executa a operação e retorna as primeiras N linhas |
ResultFrame |
|
|
Executa a operação e retorna as últimas N linhas |
ResultFrame |
|
|
Salva os resultados em uma tabela do MaxCompute e retorna um novo DataFrame apontando para essa tabela |
PyODPS DataFrame |
|
|
Converte uma Collection em um pandas DataFrame ou uma Sequence em uma Series. Defina |
pandas DataFrame ou PyODPS DataFrame |
|
|
Métodos de plotagem |
N/A |
Em ambientes interativos (como notebooks Jupyter), o PyODPS DataFrame chama automaticamente execute ao exibir resultados ou chamar repr. Nenhuma chamada manual é necessária.
# Non-interactive environment: call execute() explicitly
print(iris[iris.sepallength < 5][:5].execute())
# Interactive environment: execute is called automatically
iris[iris.sepallength < 5][:5]
Ambos produzem:
sepallength sepalwidth petallength petalwidth name
0 4.9 3.0 1.4 0.2 Iris-setosa
1 4.7 3.2 1.3 0.2 Iris-setosa
2 4.6 3.1 1.5 0.2 Iris-setosa
3 4.6 3.4 1.4 0.3 Iris-setosa
4 4.4 2.9 1.4 0.2 Iris-setosa
Para desativar a execução automática em um ambiente interativo:
from odps import options
options.interactive = False
# Now repr() displays the abstract syntax tree (AST), not results
iris[iris.sepallength < 5][:5]
Saída:
Collection: ref_0
odps.Table
name: hudi_mc_0612.`iris3`
schema:
sepallength : double # Sepal length (cm)
sepalwidth : double # Sepal width (cm)
petallength : double # Petal length (cm)
petalwidth : double # Petal width (cm)
name : string # Type
Collection: ref_1
Filter[collection]
collection: ref_0
predicate:
Less[sequence(boolean)]
sepallength = Column[sequence(float64)] 'sepallength' from collection ref_0
Scalar[int8]
5
Slice[collection]
collection: ref_1
stop:
Scalar[int8]
5
Após desativar a execução automática, chame execute explicitamente para obter os resultados.
Recuperar resultados de um ResultFrame
execute e head retornam um ResultFrame, que é um conjunto de resultados somente leitura. Use-o para iterar sobre registros ou converter para pandas. Observe que não é possível usar um ResultFrame em cálculos subsequentes do DataFrame.
Itere sobre os registros:
result = iris.head(3)
for r in result:
print(list(r))
Saída:
[4.9, 3.0, 1.4, 0.2, 'Iris-setosa']
[4.7, 3.2, 1.3, 0.2, 'Iris-setosa']
[4.6, 3.1, 1.5, 0.2, 'Iris-setosa']
Se o pandas estiver instalado, converta um ResultFrame em um pandas DataFrame ou em um PyODPS DataFrame:
# Returns a pandas DataFrame
pd_df = iris.head(3).to_pandas()
# Returns a PyODPS DataFrame backed by pandas
wrapped_df = iris.head(3).to_pandas(wrap=True)
Como alternativa, useopen_readercomreader.to_pandas()para converter resultados em um pandas DataFrame. Consulte Tabelas .
Salvar resultados em tabelas do MaxCompute
Use persist para gravar resultados no MaxCompute e continuar trabalhando com a saída como um DataFrame. Diferentemente de execute, que retorna um ResultFrame para uso local, persist grava em uma tabela e retorna um novo PyODPS DataFrame apontando para ela.
Salvar em uma nova tabela
Passe o nome da tabela para persist:
iris2 = iris[iris.sepalwidth < 2.5].persist('pyodps_iris')
print(iris2.head(5))
Saída:
sepallength sepalwidth petallength petalwidth name
0 4.5 2.3 1.3 0.3 Iris-setosa
1 5.5 2.3 4.0 1.3 Iris-versicolor
2 4.9 2.4 3.3 1.0 Iris-versicolor
3 5.0 2.0 3.5 1.0 Iris-versicolor
4 6.0 2.2 4.0 1.0 Iris-versicolor
Salvar em uma tabela particionada
Use o parâmetro partitions para criar uma tabela particionada. A tabela será particionada nas colunas especificadas:
iris3 = iris[iris.sepalwidth < 2.5].persist('pyodps_iris_test', partitions=['name'])
print(iris3.data)
Saída:
odps.Table
name: odps_test_sqltask_finance.`pyodps_iris`
schema:
sepallength : double
sepalwidth : double
petallength : double
petalwidth : double
partitions:
name : string
Gravar em uma partição existente
Use o parâmetro partition para direcionar a gravação a uma partição específica de uma tabela existente (por exemplo, ds=test). A tabela deve conter todas as colunas do DataFrame com tipos correspondentes.
drop_partition=True: remove a partição se ela já existir.create_partition=True: cria a partição se ela não existir.
Esses parâmetros são válidos apenas quando partition for especificado.
print(iris[iris.sepalwidth < 2.5].persist(
'pyodps_iris_partition',
partition='ds=test',
drop_partition=True,
create_partition=True
).head(5))
Saída:
sepallength sepalwidth petallength petalwidth name ds
0 4.5 2.3 1.3 0.3 Iris-setosa test
1 5.5 2.3 4.0 1.3 Iris-versicolor test
2 4.9 2.4 3.3 1.0 Iris-versicolor test
3 5.0 2.0 3.5 1.0 Iris-versicolor test
4 6.0 2.2 4.0 1.0 Iris-versicolor test
Definir um ciclo de vida (time-to-live)
Use o parâmetro lifecycle para definir por quantos dias os dados da tabela serão retidos. Por exemplo, defina um ciclo de vida de 10 dias:
print(iris[iris.sepalwidth < 2.5].persist('pyodps_iris', lifecycle=10).head(5))
Saída:
sepallength sepalwidth petallength petalwidth name
0 4.5 2.3 1.3 0.3 Iris-setosa
1 5.5 2.3 4.0 1.3 Iris-versicolor
2 4.9 2.4 3.3 1.0 Iris-versicolor
3 5.0 2.0 3.5 1.0 Iris-versicolor
4 6.0 2.2 4.0 1.0 Iris-versicolor
Persistir a partir de uma fonte de dados exclusiva do pandas
Se o seu DataFrame não contiver objetos do MaxCompute (apenas objetos do pandas), especifique o objeto de entrada do MaxCompute ao chamar persist:
# Option 1: pass the entrance object directly
df.persist('table_name', odps=o)
# Option 2: mark the entrance object as global
o.to_global()
df.persist('table_name')
Salvar resultados em um pandas DataFrame
Chame to_pandas para converter resultados em um pandas DataFrame. Defina wrap=True para obter um PyODPS DataFrame.
# Returns a pandas DataFrame
print(type(iris[iris.sepalwidth < 2.5].to_pandas()))
# <class 'pandas.core.frame.DataFrame'>
# Returns a PyODPS DataFrame
print(type(iris[iris.sepalwidth < 2.5].to_pandas(wrap=True)))
# <class 'odps.df.core.DataFrame'>
Configurar parâmetros de tempo de execução
Passe o parâmetro hints para execute, persist ou to_pandas para definir parâmetros de tempo de execução específicos para essa chamada. Isso funciona apenas com o backend SQL do MaxCompute.
print(iris[iris.sepallength < 5].to_pandas(hints={'odps.sql.mapper.split.size': 16}))
Saída:
sepallength sepalwidth petallength petalwidth name
0 4.5 2.3 1.3 0.3 Iris-setosa
1 4.9 2.4 3.3 1.0 Iris-versicolor
Para definir parâmetros globais de tempo de execução, consulte SQL.
Visualizar detalhes de tempo de execução
Defina options.verbose = True para imprimir o SQL compilado, o ID da instância e a URL do LogView para cada operação:
from odps import options
options.verbose = True
print(iris[iris.sepallength < 5].exclude('sepallength')[:5].execute())
Saída:
Sql compiled:
SELECT t1.`sepalwidth`, t1.`petallength`, t1.`petalwidth`, t1.`name`
FROM odps_test_sqltask_finance.`pyodps_iris` t1
WHERE t1.`sepallength` < 5
LIMIT 5
Instance ID:
Log view:http://logview
sepalwidth petallength petalwidth name
0 2.3 1.3 0.3 Iris-setosa
1 2.4 3.3 1.0 Iris-versicolor
Para capturar a saída de log no código, atribua uma função personalizada a options.verbose_log:
my_logs = []
def my_logger(x):
my_logs.append(x)
options.verbose_log = my_logger
print(iris[iris.sepallength < 5].exclude('sepallength')[:5].execute())
print(my_logs)
Saída:
sepalwidth petallength petalwidth name
0 2.3 1.3 0.3 Iris-setosa
1 2.4 3.3 1.0 Iris-versicolor
['Sql compiled:', 'CREATE TABLE tmp_pyodps_24332bdb_4fd0_4d0d_aed4_38a443618268 LIFECYCLE 1 AS \nSELECT t1.`sepalwidth`, t1.`petallength`, t1.`petalwidth`, t1.`name` \nFROM odps_test_sqltask_finance.`pyodps_iris` t1 \nWHERE t1.`sepallength` < 5 \nLIMIT 5', 'Instance ID: 20230815034706122gbymevg*****', ' Log view:]
Armazenar resultados intermediários em cache
Quando várias operações subsequentes compartilham a mesma Collection intermediária custosa, use cache para evitar recálculos repetidos. Marque a Collection com cache antes da ramificação. A execução continua adiada e não começa no momento da chamada de cache.
cached = iris[iris.sepalwidth < 3.5]['sepallength', 'name'].cache()
df = cached.head(3)
print(df)
# The following result is returned:
sepallength name
0 4.5 Iris-setosa
1 5.5 Iris-versicolor
2 4.9 Iris-versicolor
# cached is already computed, so this returns immediately without re-executing
print(cached.head(3))
# The following result is returned:
sepallength name
0 4.5 Iris-setosa
1 5.5 Iris-versicolor
2 4.9 Iris-versicolor
Executar operações de forma assíncrona e em paralelo
Execução assíncrona
Passe async_=True para execute, persist, head, tail ou to_pandas para executar a operação de forma assíncrona. O método retorna imediatamente um objeto Future. Use o parâmetro timeout para definir um tempo limite.
future = iris[iris.sepalwidth < 10].head(10, async_=True)
print(future.result())
Saída:
sepallength sepalwidth petallength petalwidth name
0 4.5 2.3 1.3 0.3 Iris-setosa
1 5.5 2.3 4.0 1.3 Iris-versicolor
2 4.9 2.4 3.3 1.0 Iris-versicolor
3 5.0 2.0 3.5 1.0 Iris-versicolor
4 6.0 2.2 4.0 1.0 Iris-versicolor
5 6.2 2.2 4.5 1.5 Iris-versicolor
6 5.5 2.4 3.8 1.1 Iris-versicolor
7 5.5 2.4 3.7 1.0 Iris-versicolor
8 6.3 2.3 4.4 1.3 Iris-versicolor
9 5.0 2.3 3.3 1.0 Iris-versicolor
Execução paralela com a API Delay
Use a API Delay para adiar chamadas de execute, persist, head, tail e to_pandas. Ao chamar delay.execute, o sistema identifica dependências entre as operações adiadas e as executa com base na concorrência especificada. Os métodos adiados retornam objetos Future. A execução só começa quando delay.execute é chamado.
from odps.df import Delay
delay = Delay() # Create a Delay object
df = iris[iris.sepal_width < 5].cache() # Shared base for all branches
# Register each branch — returns Future objects; execution has not started yet
future1 = df.sepal_width.sum().execute(delay=delay)
future2 = df.sepal_width.mean().execute(delay=delay)
future3 = df.sepal_length.max().execute(delay=delay)
# Start execution with 3 concurrent threads
delay.execute(n_parallel=3)
# |==========================================| 1 / 1 (100.00%) 21s
print(future1.result())
# 25.0
print(future2.result())
# 2.272727272727273
O PyODPS DataFrame executa primeiro o objeto compartilhado df e depois roda de future1 a future3 com a concorrência especificada.
Passe async_=True para delay.execute para executar todo o lote de forma assíncrona. Use o parâmetro timeout para definir um tempo limite para o lote.