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 |
|
|
callable |
Une UDF qui accepte un pandas DataFrame (un lot de lignes) et renvoie un pandas DataFrame ou une Series. |
|
|
int |
Nombre maximal de lignes par lot. Ce paramètre contrôle l'utilisation de la mémoire et le degré de parallélisme. |
|
|
str |
Type de sortie : |
|
|
pd.Series |
Types de données des colonnes de sortie. Ils doivent correspondre exactement aux colonnes renvoyées par |
|
|
Index |
Objet index de sortie. |
|
|
IndexValue |
Métadonnées de l'index distribué. Récupérez cette valeur depuis le DataFrame d'origine. |
|
|
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()
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 |
|
|
Nom du projet MaxCompute |
|
|
|
ID de la région |
|
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.