Ao executar df.apply() em um DataFrame distribuído do MaxFrame, o processamento linha a linha gera gargalos em grande escala. O operador apply_chunk resolve esse problema ao processar dados em lotes configuráveis e enviar cada lote a um worker do MaxCompute como um DataFrame do pandas. Para cargas de trabalho sensíveis a desempenho, use apply_chunk em vez de df.apply().
Use apply_chunk quando:
Sua função definida pelo usuário (UDF) operar em um DataFrame do pandas e precisar ser executada em escala
Você desejar controlar o uso de memória e o grau de paralelismo
O processamento linha a linha com
df.apply()for muito lento para o seu volume de dados
Como funciona
O operador apply_chunk divide o DataFrame distribuído em partes, envia cada parte a um worker do MaxCompute como um DataFrame do pandas e mescla os resultados. Os parâmetros principais — batch_rows, output_type e dtypes — indicam ao MaxFrame como particionar os dados, o que a UDF retorna e como validar e mesclar a saída. Também é possível usar decoradores de UDF, como @with_python_requirements, para gerenciar dependências em tarefas complexas.
Parâmetros
DataFrame.mf.apply_chunk(
func,
batch_rows=None,
output_type=None,
dtypes=None,
index=None,
index_value=None,
columns=None,
elementwise=None,
sort=False,
**kwds
)
|
Parâmetro |
Tipo |
Descrição |
|
|
callable |
UDF que aceita um DataFrame do pandas (um lote de linhas) e retorna um DataFrame ou Series do pandas. |
|
|
int |
Número máximo de linhas por lote. Controla o uso de memória e o grau de paralelismo. |
|
|
str |
Tipo de saída: |
|
|
pd.Series |
Tipos de dados das colunas de saída. Devem corresponder exatamente às colunas retornadas por |
|
|
Index |
Objeto de índice de saída. |
|
|
IndexValue |
Metadados do índice distribuído. Obtenha este valor do DataFrame original. |
|
|
bool |
Indica se os dados devem ser ordenados dentro dos grupos em um cenário de groupby. |
Entendendo dtypes
O parâmetro dtypes é um objeto pd.Series que descreve os nomes das colunas e os tipos de dados da saída da sua UDF. O MaxFrame utiliza essa informação para validar e mesclar resultados entre workers. Se dtypes não corresponder ao que func realmente retorna, o job falhará durante a execução.
Para inspecionar a estrutura de dtypes de um DataFrame:
>>> df.dtypes
A object
B object
dtype: object
A forma mais comum de fornecer dtypes é copiá-lo do DataFrame original:
dtypes=df.dtypes.copy()
Nunca passe df.dtypes diretamente. Sempre use .copy() para evitar a modificação dos metadados do DataFrame original.
Caso sua UDF altere o esquema de saída (adicionando ou removendo colunas), construa um novo objeto pd.Series que descreva a saída real:
import pandas as pd
# UDF drops column 'A' and adds column 'C'
new_dtypes = pd.Series({
"B": df.dtypes["B"],
"C": pd.StringDtype(),
})
Exemplo
O exemplo abaixo cria um DataFrame do MaxFrame com uma coluna do tipo dict e usa apply_chunk para atualizar valores em cada lote.
import os
import pyarrow as pa
import pandas as pd
import maxframe.dataframe as md
from maxframe.lib.dtypes_extension import dict_
from maxframe import new_session
from odps import ODPS
o = ODPS(
# Get credentials from environment variables.
# Never hardcode AccessKey ID or AccessKey secret in your code.
os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
project='<your-project>',
endpoint='https://service.cn-<your-region>.maxcompute.aliyun.com/api',
)
session = new_session(o)
# Create a MaxFrame DataFrame with a dict-type column
col_a = pd.Series(
data=[[("k1", 1), ("k2", 2)], [("k1", 3)], None],
index=[1, 2, 3],
dtype=dict_(pa.string(), pa.int64()),
)
col_b = pd.Series(
data=["A", "B", "C"],
index=[1, 2, 3],
)
df = md.DataFrame({"A": col_a, "B": col_b})
df.execute()
def custom_set_item(df: pd.DataFrame) -> pd.DataFrame:
"""Add key 'x' with value 100 to each non-null dict in column A."""
for name, value in df["A"].items():
if value is not None:
df["A"][name]["x"] = 100
return df
result_df = df.mf.apply_chunk(
custom_set_item,
output_type="dataframe",
dtypes=df.dtypes.copy(), # Must match the columns returned by the UDF
batch_rows=2, # Process 2 rows per chunk
skip_infer=True,
index=df.index,
).execute()
session.destroy()
Substitua os placeholders a seguir pelos valores reais:
|
Placeholder |
Descrição |
Exemplo |
|
|
Nome do projeto MaxCompute |
|
|
|
ID da região |
|
Ajuste de desempenho
Defina batch_rows com base nos seus dados e recursos
O parâmetro batch_rows controla quantas linhas cada worker processa por lote. O valor ideal depende do tamanho da linha e da memória disponível:
|
Direção |
Efeito |
Risco |
|
Valores maiores |
Menos lotes, menor sobrecarga de agendamento, melhor throughput |
Erros de falta de memória (OOM) se as linhas forem largas ou a UDF consumir muita memória |
|
Valores menores |
Mais lotes, maior grau de paralelismo |
Maior sobrecarga de agendamento para lotes muito pequenos |
Comece com um valor conservador, monitore o uso de memória no LogView e faça ajustes a partir daí. Se ocorrerem erros de OOM, reduza batch_rows. Caso os jobs sejam executados rapidamente com baixo uso de memória, aumente o valor.
Se sua UDF exigir mais memória por lote, utilize o decorador @with_running_options para elevar o limite de memória por tarefa:
from maxframe.udf import with_running_options
@with_running_options(memory=16)
def my_udf(df: pd.DataFrame) -> pd.DataFrame:
...
Sempre declare output_type e dtypes explicitamente
Por padrão, o MaxFrame tenta inferir o esquema de saída executando sua função em dados de amostra. A inferência pode falhar ou produzir resultados incorretos para esquemas complexos. Declarar explicitamente output_type e dtypes evita falhas em tempo de execução e acelera o job.
Sem declaração explícita (falha em tempo de execução):
result_df = df.mf.apply_chunk(process) # dtypes is missing
Com declaração explícita (correto):
result_df = df.mf.apply_chunk(
process,
output_type="dataframe",
dtypes=df.dtypes.copy(),
)
Um erro comum é passar dtypes que não correspondem à saída real da UDF. Por exemplo, se você remover uma coluna dentro de func, mas ainda passar o df.dtypes original, o MaxFrame lançará um ValueError porque o esquema declarado terá mais colunas do que o DataFrame retornado.
Retorne apenas as colunas necessárias
Cada coluna de saída é serializada, transferida e mesclada entre workers. Retornar colunas não utilizadas aumenta o uso de memória e torna o job mais lento. Dentro de func, remova quaisquer colunas desnecessárias antes de retornar o resultado.
Depure UDFs com print e flush=True
As UDFs são executadas em workers remotos do MaxCompute. Use print(..., flush=True) para emitir logs imediatamente — eles aparecerão no LogView. Envolva o corpo da UDF em um bloco try/except para capturar e registrar erros antes que eles surjam como falhas genéricas:
def process(chunk: pd.DataFrame) -> pd.DataFrame:
try:
print(f"Processing chunk: shape={chunk.shape}, columns={list(chunk.columns)}", flush=True)
result = chunk.sort_values("B")
print("Chunk processed successfully.", flush=True)
return result
except Exception as e:
print(f"[ERROR] {type(e).__name__}: {e}", flush=True)
raise
Perguntas frequentes
Por que recebo TypeError: cannot determine dtype?
O MaxFrame não conseguiu inferir o esquema de saída. Passe dtypes e output_type explicitamente ao chamar apply_chunk.
Por que minha saída está vazia ou com colunas ausentes?
O dtypes fornecido não corresponde ao retorno da sua UDF. Verifique se os nomes das colunas em dtypes correspondem exatamente às colunas no DataFrame retornado por func. Remova de dtypes qualquer coluna que func não inclua em seu valor de retorno.
Por que o job trava ou atinge o tempo limite?
Provavelmente batch_rows está muito alto. Reduza esse valor e aloque mais recursos para dar a cada worker um lote menor e mais rápido de processar. Verifique também se há erros de OOM no LogView — eles frequentemente aparecem como travamentos em vez de falhas explícitas.