Montez et utilisez Alibaba Cloud OSS en tant que stockage distribué dans MaxFrame à l'aide du décorateur with_fs_mount. FS Mount offre un accès stable, au niveau du système de fichiers, aux données externes pour le traitement de données à grande échelle.
Cas d'utilisation
FS Mount est idéal pour l'analyse de big data dans les jobs MaxFrame qui interagissent avec un stockage objet persistant tel qu'OSS. Par exemple :
Chargez, nettoyez et traitez des données brutes depuis OSS.
Écrivez des résultats intermédiaires dans OSS pour les tâches en aval.
Partagez des ressources statiques, telles que des fichiers de modèle entraînés et des fichiers de configuration.
Les méthodes traditionnelles de lecture/écriture telles que pd.read_csv("oss://...") sont limitées par les performances du SDK et la surcharge réseau dans un environnement distribué. FS Mount vous permet d'accéder aux fichiers OSS comme s'ils se trouvaient sur un disque local, ce qui améliore considérablement l'efficacité du développement.
Procédure
Activation des services et attribution des autorisations
-
Activer OSS et créer un bucket.
Connectez-vous à la console OSS.
Dans le volet de navigation de gauche, cliquez sur Buckets.
-
Sur la page Buckets, cliquez sur Create Bucket.
Dans cet exemple, le nom du bucket est
xxx-oss-test-sh.
-
Créez un rôle RAM pour MaxCompute et accordez-lui l'accès à l'environnement d'exécution.
Connectez-vous à la console RAM.
Dans la barre de navigation de gauche, sélectionnez .
Sur la page Roles, cliquez sur Create Role.
-
Dans le coin supérieur droit de la page Create Role, cliquez sur Create Service Linked Role.
Sur la page Create Role, définissez Principal Type sur Cloud Service.
Pour Principal Type, sélectionnez MaxCompute.
-
Sous l'onglet Manage Permissions, cliquez sur Create Authorization. Dans le panneau Create Authorization qui s'affiche, sélectionnez les stratégies à attribuer au rôle, puis cliquez sur OK.
Sélectionnez les stratégies suivantes :
AliyunOSSFullAccess : accorde des autorisations de gestion d'OSS.
AliyunMaxComputeFullAccess : accorde des autorisations de gestion de MaxCompute.
Utilisation de with_fs_mount pour monter OSS
-
Recommandé : authentification avec un ARN de rôle
from maxframe.udf import with_fs_mount @with_fs_mount( "oss://oss-cn-xxxx-internal.aliyuncs.com/xxx-oss-test-sh/test/", "/mnt/oss_data", storage_options={ "role_arn": "acs:ram::xxx:role/maxframe-oss" }, ) def _process(batch_df): import os if os.path.exists('/mnt/oss_data'): print(f"Mounted files: {os.listdir('/mnt/oss_data')}") else: print("/mnt/oss_data not mounted!") return batch_df * 2 -
Non recommandé : identifiants codés en dur
Cette méthode est réservée aux tests et n'est pas recommandée pour les environnements de production.
storage_options={ "access_key_id": "LTAI5t...", "access_key_secret": "Wp9H..." }ImportantÉvitez de coder en dur votre AccessKey. Utilisez
role_arnpour permettre au système de demander automatiquement un jeton STS temporaire, empêchant ainsi la divulgation de votre paire AccessKey.
Utilisation de with_running_options pour contrôler l'allocation des ressources
Utilisez le décorateur with_running_options pour allouer des ressources CPU et mémoire à votre tâche :
from maxframe.udf import with_running_options
@with_running_options(engine="dpe", cpu=2, memory=16)
@with_fs_mount(...)
def _process(batch_df):
...
|
Paramètre |
Valeur recommandée |
Description |
|
|
Fixe |
FS Mount prend actuellement en charge uniquement le moteur DPE. |
|
|
1–4 |
Augmentez cette valeur pour les tâches intensives en E/S ou en décompression. |
|
|
Commencez par 8 Go |
Pour le chargement de fichiers volumineux, 16 Go ou plus sont recommandés. |
Exemple
Modèle recommandé : traitez les données par lots.
Dans les scénarios de traitement de données à grande échelle, utilisez la fonctionnalité apply_chunk de MaxFrame pour traiter les données d'entrée par lots.
Création d'une session MaxFrame
import os
from odps import ODPS
from maxframe import new_session
from maxframe.udf import with_fs_mount, with_running_options
# Initialize the ODPS client.
# We recommend setting the ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET
# environment variables instead of hard-coding the AccessKey ID and AccessKey Secret strings.
o = ODPS(
os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
project='<your-project>',
endpoint='https://service.cn-<region>.maxcompute.aliyun.com/api',
)
# Set the runtime image.
# The `maxframe_service_dpe_runtime` image includes the required ossfs2 dependency.
# If you use a custom image, you must download the dependency and include it in your image.
# You can find the package at the link below this code block.
options.sql.settings = { "odps.session.image": "maxframe_service_dpe_runtime"}
# Start the session.
session = new_session(o)
print("LogView:", session.get_logview_address())
print("Session ID:", session.session_id)
@with_running_options(engine="dpe", cpu=2, memory=8)
@with_fs_mount(
"oss://oss-cn-<region>-internal.aliyuncs.com/wzy-oss-test-sh/test/",
"/mnt/oss_data",
storage_options={
"role_arn": "acs:ram::<uid>:role/maxframe-oss"
},
)
Package de dépendance OSSFS : ossfs2_2.0.3.1_linux_x86_64.deb
Création d'une fonction définie par l'utilisateur (UDF)
def _process(batch_df):
import pandas as pd
import os
# Step 1: Check if the mount was successful.
mount_point = "/mnt/oss_data"
if not os.path.exists(mount_point):
raise RuntimeError("OSS mount failed!")
# Step 2: Load data, such as a mapping table or dictionary.
mapping_file = os.path.join(mount_point, "category_map.csv")
if os.path.isfile(mapping_file):
mapping_df = pd.read_csv(mapping_file)
# Step 3: Process the current chunk.
result = batch_df.copy()
result['F'] = result['A'] * 10
return result
Construction d'un DataFrame et application de l'UDF
import maxframe.dataframe as md
data = [[1.0, 2.0, 3.0, 4.0, 5.0], ...]
df = md.DataFrame(data, columns=['A', 'B', 'C', 'D', 'E'])
# Apply the UDF by using apply_chunk.
result_df = df.mf.apply_chunk(
_process,
skip_infer=True,
output_type="dataframe",
dtypes=df.dtypes,
index=df.index
)
# Execute the operation and fetch the result.
result = result_df.execute().fetch()
La définition de skip_infer=True ignore l'inférence de type et améliore la vitesse d'exécution. Toutefois, assurez-vous que dtypes et index sont transmis correctement.
Dépannage
Vérification de l'état du montage
Ajoutez des journaux de débogage dans la fonction _process :
import os
print("Mount path exists:", os.path.exists("/mnt/oss_data"))
print("Files in mount:", os.listdir("/mnt/oss_data") if os.path.exists("/mnt/oss_data") else [])
Consultez la sortie LogView pour confirmer que des journaux similaires aux suivants sont générés :
FS Mount successful! /mnt/oss_data: ['data.csv', 'config.json', 'model.pkl']
Processing batch with shape: (1000, 5)