Tous les produits
Search
Centre de documentation

MaxCompute:Bonnes pratiques pour l'utilisation de l'opérateur MaxFrame apply_chunk

Dernière mise à jour :Aug 10, 2026

Lorsque vous exécutez df.apply() sur un DataFrame MaxFrame distribué, le traitement ligne par ligne des données crée un goulot d'étranglement à grande échelle. L'opérateur apply_chunk résout ce problème en traitant les données par lots configurables et en envoyant chaque lot à un worker MaxCompute sous la forme d'un pandas DataFrame. Pour les charges de travail sensibles aux performances, privilégiez apply_chunk plutôt que df.apply().

Utilisez apply_chunk dans les cas suivants :

  • Votre fonction définie par l'utilisateur (UDF) opère sur un pandas DataFrame et doit s'exécuter à grande échelle.

  • Vous souhaitez contrôler l'utilisation de la mémoire et le degré de parallélisme.

  • Le traitement ligne par ligne avec df.apply() est trop lent pour votre volume de données.

Fonctionnement

L'opérateur apply_chunk divise le DataFrame distribué en lots, envoie chaque lot à un worker MaxCompute sous forme de pandas DataFrame, puis fusionne les résultats. Les paramètres clés — batch_rows, output_type et dtypes — indiquent à MaxFrame comment partitionner les données, quel type de résultat l'UDF renvoie, ainsi que la manière de valider et de fusionner la sortie. Vous pouvez également utiliser des décorateurs UDF tels que @with_python_requirements pour gérer les dépendances liées aux tâches complexes.

Paramètres

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
)

Paramètre

Type

Description

func

callable

Une UDF qui accepte un pandas DataFrame (un lot de lignes) et renvoie un pandas DataFrame ou une Series.

batch_rows

int

Nombre maximal de lignes par lot. Ce paramètre contrôle l'utilisation de la mémoire et le degré de parallélisme.

output_type

str

Type de sortie : "dataframe" ou "series".

dtypes

pd.Series

Types de données des colonnes de sortie. Ils doivent correspondre exactement aux colonnes renvoyées par func.

index

Index

Objet index de sortie.

index_value

IndexValue

Métadonnées de l'index distribué. Récupérez cette valeur depuis le DataFrame d'origine.

sort

bool

Indique s'il faut trier les données au sein des groupes dans un scénario groupby.

Compréhension du paramètre dtypes

Le paramètre dtypes est un objet pd.Series qui décrit les noms et les types de données des colonnes de la sortie de votre UDF. MaxFrame l'utilise pour valider et fusionner les résultats entre les workers. Si la valeur de dtypes ne correspond pas à ce que func renvoie réellement, le job échoue lors de l'exécution.

Pour inspecter la structure de dtypes pour un DataFrame :

>>> df.dtypes
A    object
B    object
dtype: object

La méthode la plus courante pour fournir dtypes consiste à le copier depuis le DataFrame d'origine :

dtypes=df.dtypes.copy()
Important

Ne transmettez jamais directement df.dtypes. Utilisez toujours la méthode .copy() afin d'éviter toute modification des métadonnées du DataFrame d'origine.

Si votre UDF modifie le schéma de sortie (ajout ou suppression de colonnes), construisez un nouvel objet pd.Series décrivant la sortie réelle :

import pandas as pd

# UDF drops column 'A' and adds column 'C'
new_dtypes = pd.Series({
    "B": df.dtypes["B"],
    "C": pd.StringDtype(),
})

Exemple

L'exemple suivant crée un DataFrame MaxFrame avec une colonne de type dict, puis utilise apply_chunk pour mettre à jour les valeurs de chaque lot.

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()

Remplacez les espaces réservés suivants par les valeurs réelles :

Espace réservé

Description

Exemple

<your-project>

Nom du projet MaxCompute

my_project

<your-region>

ID de la région

cn-hangzhou

Optimisation des performances

Définissez batch_rows en fonction de vos données et ressources

Le paramètre batch_rows contrôle le nombre de lignes traitées par chaque worker pour un lot donné. La valeur appropriée dépend de la taille de vos lignes et de la mémoire disponible :

Orientation

Effet

Risque

Valeurs plus élevées

Moins de lots, moins de frais généraux de planification, meilleur débit

Erreurs d'épuisement de la mémoire (OOM) si les lignes sont larges ou si l'UDF consomme beaucoup de mémoire

Valeurs plus faibles

Plus de lots, degré de parallélisme plus élevé

Frais généraux de planification accrus pour des lots très petits

Commencez par une valeur conservative, surveillez l'utilisation de la mémoire dans LogView et ajustez-la en conséquence. Si vous observez des erreurs OOM, réduisez la valeur de batch_rows. Si les jobs s'exécutent rapidement avec une faible utilisation de la mémoire, augmentez-la.

Si votre UDF nécessite plus de mémoire par lot, utilisez le décorateur @with_running_options pour augmenter la limite de mémoire par tâche :

from maxframe.udf import with_running_options

@with_running_options(memory=16)
def my_udf(df: pd.DataFrame) -> pd.DataFrame:
    ...

Déclarez toujours explicitement output_type et dtypes

Par défaut, MaxFrame tente de déduire le schéma de sortie en exécutant votre fonction sur un échantillon de données. Cette inférence peut échouer ou produire des résultats incorrects pour des schémas complexes. La déclaration explicite de output_type et de dtypes permet d'éviter les échecs d'exécution et d'accélérer le job.

Sans déclaration explicite (échec à l'exécution) :

result_df = df.mf.apply_chunk(process)  # dtypes is missing

Avec déclaration explicite (correct) :

result_df = df.mf.apply_chunk(
    process,
    output_type="dataframe",
    dtypes=df.dtypes.copy(),
)

Une erreur fréquente consiste à transmettre une valeur dtypes qui ne correspond pas à la sortie réelle de l'UDF. Par exemple, si vous supprimez une colonne à l'intérieur de func tout en transmettant l'objet df.dtypes d'origine, MaxFrame lève une exception ValueError car le schéma déclaré contient plus de colonnes que le DataFrame renvoyé.

Ne renvoyez que les colonnes nécessaires

Chaque colonne de sortie est sérialisée, transférée et fusionnée entre les workers. Le renvoi de colonnes inutilisées augmente l'utilisation de la mémoire et ralentit le job. À l'intérieur de func, supprimez toutes les colonnes dont vous n'avez pas besoin avant de retourner le résultat.

Débogage des UDF avec print et flush=True

Les UDF s'exécutent sur des workers MaxCompute distants. Utilisez print(..., flush=True) pour émettre les journaux immédiatement ; ils apparaissent dans LogView. Encapsulez le corps de l'UDF dans un bloc try/except afin d'intercepter et de journaliser les erreurs avant qu'elles ne se manifestent sous forme d'échecs génériques :

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

FAQ

Pourquoi obtiens-je l'erreur TypeError: cannot determine dtype ?

MaxFrame n'a pas réussi à déduire le schéma de sortie. Transmettez explicitement les paramètres dtypes et output_type lors de l'appel à apply_chunk.

Pourquoi ma sortie est-elle vide ou manquante de colonnes ?

La valeur dtypes transmise ne correspond pas à ce que renvoie votre UDF. Vérifiez que les noms de colonnes dans dtypes correspondent exactement aux colonnes du DataFrame renvoyé par func. Supprimez de dtypes toute colonne que func n'inclut pas dans sa valeur de retour.

Pourquoi le job reste-t-il bloqué ou expire-t-il ?

La valeur de batch_rows est probablement trop élevée. Réduisez-la et allouez davantage de ressources afin que chaque worker traite un lot plus petit et plus rapide. Vérifiez également LogView pour détecter d'éventuelles erreurs OOM, qui se manifestent souvent par des blocages plutôt que par des échecs explicites.