Utilitaires Microsoft Spark (MSSparkUtils) pour Fabric

Microsoft Spark Utilities (MSSparkUtils) est un paquet intégré qui vous aide à effectuer facilement des tâches courantes. Utilisez MSSparkUtils pour travailler avec des systèmes de fichiers, récupérer des variables d’environnement, enchaîner des notebooks et gérer les secrets. Le package MSSparkUtils est disponible dans les notebooks PySpark (Python), Scala et SparkR, ainsi que dans les pipelines Fabric.

Remarque

  • MsSparkUtils est officiellement renommé NotebookUtils. Le code existant restera rétrocompatible et n'entraînera aucune modification radicale. Nous recommandons vivement de passer à notebookutils afin d’assurer un support continu et un accès aux nouvelles fonctionnalités. L’espace de noms mssparkutils sera supprimé à l’avenir.
  • NotebookUtils est conçu pour fonctionner avec Spark 3.4(Runtime v1.2) et les versions ultérieures. Toutes les nouvelles fonctionnalités et mises à jour sont exclusivement supportées par l’espace de noms notebookutils à l’avenir.

Utilitaires du système de fichiers

mssparkutils.fs fournit des utilitaires pour travailler avec divers systèmes de fichiers, notamment Azure Data Lake Storage Gen2 et Stockage Blob Azure. Veillez à configurer l’accès à Azure Data Lake Storage Gen2 et Stockage Blob Azure de manière appropriée.

Exécutez les commandes suivantes pour obtenir une vue d’ensemble des méthodes disponibles :

from notebookutils import mssparkutils
mssparkutils.fs.help()

Sortie

mssparkutils.fs provides utilities for working with various FileSystems.

Below is overview about the available methods:

cp(from: String, to: String, recurse: Boolean = false): Boolean -> Copies a file or directory, possibly across FileSystems
mv(from: String, to: String, recurse: Boolean = false): Boolean -> Moves a file or directory, possibly across FileSystems
ls(dir: String): Array -> Lists the contents of a directory
mkdirs(dir: String): Boolean -> Creates the given directory if it does not exist, also creating any necessary parent directories
put(file: String, contents: String, overwrite: Boolean = false): Boolean -> Writes the given String out to a file, encoded in UTF-8
head(file: String, maxBytes: int = 1024 * 100): String -> Returns up to the first 'maxBytes' bytes of the given file as a String encoded in UTF-8
append(file: String, content: String, createFileIfNotExists: Boolean): Boolean -> Append the content to a file
rm(dir: String, recurse: Boolean = false): Boolean -> Removes a file or directory
exists(file: String): Boolean -> Check if a file or directory exists
mount(source: String, mountPoint: String, extraConfigs: Map[String, Any]): Boolean -> Mounts the given remote storage directory at the given mount point
unmount(mountPoint: String): Boolean -> Deletes a mount point
mounts(): Array[MountPointInfo] -> Show information about what is mounted
getMountPath(mountPoint: String, scope: String = ""): String -> Gets the local path of the mount point

Use mssparkutils.fs.help("methodName") for more info about a method.

MSSparkUtils fonctionne avec le système de fichiers de la même façon que les API Spark. Prenons par exemple mssparkuitls.fs.mkdirs() et l’utilisation de lakehouse :

Utilisation Chemin d’accès relatif à partir de la racine HDFS Chemin d’accès absolu pour le système de fichiers ABFS Chemin d’accès absolu pour le système de fichiers local dans le nœud de pilote
Lakehouse autre que celui par défaut Non pris en charge mssparkutils.fs.mkdirs(« abfss://< container_name>@<storage_account_name.dfs.core.windows.net/>< new_dir> ») mssparkutils.fs.mkdirs(« file :/<new_dir> »)
Lakehouse par défaut Répertoire sous « Files » ou « Tables » : mssparkutils.fs.mkdirs("Files/<new_dir>") mssparkutils.fs.mkdirs(« abfss://< container_name>@<storage_account_name.dfs.core.windows.net/>< new_dir> ») mssparkutils.fs.mkdirs(« file :/<new_dir> »)

Énumérer des fichiers

Pour répertorier le contenu d’un répertoire, utilisez mssparkutils.fs.ls(’Le chemin d’accès à votre répertoire’). Par exemple :

mssparkutils.fs.ls("Files/tmp") # works with the default lakehouse files using relative path 
mssparkutils.fs.ls("abfss://<container_name>@<storage_account_name>.dfs.core.windows.net/<path>")  # based on ABFS file system 
mssparkutils.fs.ls("file:/tmp")  # based on local file system of driver node 

Affichez les propriétés de fichier

Cette méthode renvoie les propriétés du fichier, y compris le nom du fichier, le chemin du fichier, la taille du fichier et s’il s’agit d’un répertoire ou d’un fichier.

files = mssparkutils.fs.ls('Your directory path')
for file in files:
    print(file.name, file.isDir, file.isFile, file.path, file.size)

Créer un répertoire

Cette méthode crée le répertoire spécifié s’il n’existe pas, et crée les répertoires parents nécessaires.

mssparkutils.fs.mkdirs('new directory name')  
mssparkutils.fs. mkdirs("Files/<new_dir>")  # works with the default lakehouse files using relative path 
mssparkutils.fs.ls("abfss://<container_name>@<storage_account_name>.dfs.core.windows.net/<new_dir>")  # based on ABFS file system 
mssparkutils.fs.ls("file:/<new_dir>")  # based on local file system of driver node 

Copier un fichier

Cette méthode copie un fichier ou un répertoire et prend en charge l’activité de copie entre les systèmes de fichiers.

mssparkutils.fs.cp('source file or directory', 'destination file or directory', True)# Set the third parameter as True to copy all files and directories recursively

Fichier de copie performant

Cette méthode fournit un moyen plus rapide de copier ou de déplacer des fichiers, en particulier de gros volumes de données.

mssparkutils.fs.fastcp('source file or directory', 'destination file or directory', True)# Set the third parameter as True to copy all files and directories recursively

Afficher un aperçu du contenu du fichier

Cette méthode renvoie jusqu’au premier maxBytes octet du fichier spécifié sous forme de chaîne encodée en UTF-8.

# Set the second parameter as an integer for the maxBytes to read
mssparkutils.fs.head('file path', <maxBytes>)

Déplacer un fichier

Cette méthode déplace un fichier ou un répertoire et prend en charge les déplacements entre les systèmes de fichiers.

mssparkutils.fs.mv('source file or directory', 'destination directory', True) # Set the last parameter as True to firstly create the parent directory if it does not exist
mssparkutils.fs.mv('source file or directory', 'destination directory', True, True) # Set the third parameter to True to firstly create the parent directory if it does not exist. Set the last parameter to True to overwrite the updates.

Écrire dans un fichier

Cette méthode écrit la chaîne donnée dans un fichier encodé au format UTF-8.

mssparkutils.fs.put("file path", "content to write", True) # Set the last parameter as True to overwrite the file if it existed already

Ajouter du contenu à un fichier

Cette méthode ajoute la chaîne donnée à un fichier, encodé au format UTF-8.

mssparkutils.fs.append("file path", "content to append", True) # Set the last parameter as True to create the file if it does not exist

Remarque

Lorsque vous utilisez l’API mssparkutils.fs.append en for boucle pour écrire sur le même fichier, nous vous recommandons d’ajouter une sleep instruction d’environ 0,5 à 1 seconde entre les écritures récurrentes. Le mssparkutils.fs.append fonctionnement interne flush de l’API est asynchrone, donc un court délai aide à garantir l’intégrité des données.

Supprimer un fichier ou un répertoire

Cette méthode déplace un fichier ou un répertoire.

mssparkutils.fs.rm('file path', True) # Set the last parameter as True to remove all files and directories recursively

Monter/démonter le répertoire

Pour plus d’informations sur l’utilisation détaillée, voir Montage et démontage du fichier.

Utilitaires de notebook

Utilisez les utilitaires de Notebook MSSparkUtils pour exécuter un bloc-notes ou quitter un bloc-notes avec une valeur. Exécutez la commande suivante pour obtenir une vue d’ensemble des méthodes disponibles :

mssparkutils.notebook.help()

Sortie:


exit(value: String): Raises NotebookExit Exception -> This method lets you exit a notebook with a value.
run(path: String, timeoutSeconds: int, arguments: Map): String -> This method runs a notebook and returns its exit value.

Remarque

Les utilitaires de bloc-notes ne s’appliquent pas aux définitions de travaux Apache Spark (SJD).

Référencer un cahier

Cette méthode référence un notebook et renvoie sa valeur de sortie. Vous pouvez exécuter des appels de fonction d’imbrication dans un notebook de manière interactive ou dans un pipeline. Le notebook référencé s’exécute sur le pool Spark du notebook qui appelle cette fonction.

mssparkutils.notebook.run("notebook name", <timeoutSeconds>, <parameterMap>, <workspaceId>)

Par exemple :

mssparkutils.notebook.run("Sample1", 90, {"input": 20 })

Le notebook Fabric prend également en charge le référencement de blocs-notes sur plusieurs espaces de travail en spécifiant l’ID de l’espace de travail.

mssparkutils.notebook.run("Sample1", 90, {"input": 20 }, "fe0a6e2a-a909-4aa3-a698-0a651de790aa")

Vous pouvez ouvrir le lien d’instantané de l’exécution de référence dans la sortie de cellule. L’instantané capture les résultats de l’exécution du code et vous permet de déboguer facilement une exécution de référence.

Capture d’écran montrant le résultat de l’exécution de référence.

Capture d’écran d’un instantané avec les résultats de l’exécution du code.

Remarque

  • Le notebook de référence inter-espaces de travail est pris en charge par le runtime version 1.2 et ultérieure.
  • Si vous utilisez les fichiers sous les ressources du carnet, utilisez mssparkutils.nbResPath dans le carnet référencé pour vous assurer qu’il pointe vers le même dossier que la course interactive.

Exécution de référence de plusieurs notebooks en parallèle

Important

Cette fonctionnalité est en version préliminaire.

La méthode mssparkutils.notebook.runMultiple() vous permet d’exécuter plusieurs notebooks en parallèle ou avec une structure topologique prédéfinie. L’API utilise une implémentation multithread pour soumettre, mettre en file et surveiller des carnets enfants qui s’exécutent sur des instances REPL isolées (lecture-évaluation-boucle d’impression) au sein de la session Spark existante. Les carnets enfants référencés partagent les ressources de calcul de la session.

Avec mssparkutils.notebook.runMultiple(), vous pouvez :

  • Exécutez plusieurs notebooks simultanément, sans attendre que chacun d’eux se termine.

  • Spécifiez les dépendances et l’ordre d’exécution de vos notebooks à l’aide d’un format JSON simple.

  • Optimisez l’utilisation des ressources de calcul Spark et réduisez le coût de vos projets Fabric.

  • Consultez les captures instantanées de chaque exécution de notebook dans les résultats, et déboguez et surveillez facilement vos tâches de notebook.

  • Obtenez la valeur de sortie de chaque activité exécutive et utilisez-les dans les tâches en aval.

Vous pouvez également essayer d’exécuter mssparkutils.notebook.help("runMultiple") pour rechercher l’exemple et l’utilisation détaillée.

Voici un exemple simple d’exécution d’une liste de carnets en parallèle à l’aide de cette méthode :


mssparkutils.notebook.runMultiple(["NotebookSimple", "NotebookSimple2"])

Le résultat d’exécution du notebook racine est le suivant :

Capture d’écran de la référence d’une liste de notebooks.

L’exemple suivant montre des carnets exécutants avec une structure topologique en utilisant mssparkutils.notebook.runMultiple(). Utilisez cette méthode pour orchestrer facilement les notebooks avec une expérience de code.

# run multiple notebooks with parameters
DAG = {
    "activities": [
        {
            "name": "NotebookSimple", # activity name, must be unique
            "path": "NotebookSimple", # notebook path
            "timeoutPerCellInSeconds": 90, # max timeout for each cell, default to 90 seconds
            "args": {"p1": "changed value", "p2": 100}, # notebook parameters
        },
        {
            "name": "NotebookSimple2",
            "path": "NotebookSimple2",
            "timeoutPerCellInSeconds": 120,
            "args": {"p1": "changed value 2", "p2": 200}
        },
        {
            "name": "NotebookSimple2.2",
            "path": "NotebookSimple2",
            "timeoutPerCellInSeconds": 120,
            "args": {"p1": "changed value 3", "p2": 300},
            "retry": 1,
            "retryIntervalInSeconds": 10,
            "dependencies": ["NotebookSimple"] # list of activity names that this activity depends on
        }
    ],
    "timeoutInSeconds": 43200, # max timeout for the entire DAG, default to 12 hours
    "concurrency": 50 # max number of notebooks to run concurrently, defaults to 50 but ultimately constrained by the number of driver cores
}
mssparkutils.notebook.runMultiple(DAG, {"displayDAGViaGraphviz": False})

Le résultat d’exécution du notebook racine est le suivant :

Capture d’écran de la référence d’une liste de notebooks avec des paramètres.

Remarque

  • La limite supérieure pour les activités de notebook ou les notebooks simultanés est limitée par le nombre de cœurs de pilote. Par exemple, un nœud pilote Medium avec huit cœurs peut exécuter jusqu’à huit carnets simultanément. Cette limite existe car chaque notebook soumis s’exécute sur sa propre instance REPL (read-eval-print-loop), et chaque instance consomme un cœur de pilote.
  • Le paramètre de concurrence par défaut est défini sur 50 pour prendre en charge l'ajustement automatique de la concurrence maximale lorsque les utilisateurs configurent des pools Spark avec des nœuds de plus grande taille et donc autant de cœurs de pilote. Bien que vous puissiez définir ce paramètre à une valeur plus élevée en utilisant un nœud pilote plus grand, augmenter le nombre de processus concurrents s’exécutant sur un seul nœud pilote ne s’étend généralement pas de manière linéaire. L’augmentation de la concurrence peut entraîner une réduction de l’efficacité en raison de la contention des ressources du driver et de l’exécuteur. Chaque notebook en cours d’exécution fonctionne sur une instance REPL dédiée qui consomme le CPU et la mémoire du pilote. En cas de forte concurrence, cette consommation peut augmenter le risque d’instabilité des pilotes ou d’erreurs de manque de mémoire, en particulier pour les charges de travail de longue durée.
  • Vous pourriez rencontrer des temps d’exécution plus longs pour chaque tâche individuelle à cause de la surcharge liée à l’initialisation des instances REPL et à l’orchestration de nombreux notebooks. Si des problèmes surviennent, envisagez de séparer les notebooks en plusieurs appels runMultiple ou de réduire la concurrence en ajustant le champ concurrency du paramètre DAG.
  • Lorsque vous exécutez des notebooks à courte durée de vie (par exemple, avec 5 secondes d’exécution de code), le coût d’initialisation finit par prédominer. La variabilité du temps de préparation pourrait réduire le risque de chevauchement des carnets, et donc entraîner une concurrence réalisée plus faible. Dans ces cas, il pourrait être plus optimal de combiner de petites opérations en un ou plusieurs carnets.
  • Bien que le multithreading soit utilisé pour la soumission, la mise en file d’attente et la surveillance, notez que le code exécuté dans chaque notebook n’est pas multithreadé au niveau de chaque exécuteur. Il n’y a pas de partage de ressources entre les carnets. Chaque processus de notebook se voit attribuer une partie de l’ensemble des ressources de l’exécuteur. Cette allocation peut rendre l’exécution des tâches courtes moins efficace et mettre les tâches longues en concurrence pour les ressources.
  • Le délai d’attente par défaut pour l’ensemble du DAG est de 12 heures, et le délai d’attente par défaut pour chaque cellule dans les carnets enfants est de 90 secondes. Vous pouvez modifier le délai d’expiration en définissant les champs timeoutInSeconds et timeoutPerCellInSeconds dans le paramètre DAG. À mesure que vous augmentez la concurrence, vous pourriez devoir augmenter le timeoutPerCellInSeconds pour éviter que la contention possible des ressources ne provoque des délais inutiles.

Quitter un notebook

Cette méthode quitte un notebook avec une valeur. Vous pouvez exécuter des appels de fonction d’imbrication dans un notebook de manière interactive ou dans un pipeline.

  • Lorsque vous appelez une fonction exit() à partir d’un notebook de manière interactive, le notebook Fabric lève une exception, ignore l’exécution des cellules de la sous-séquence et maintient la session Spark active.

  • Lorsque vous orchestrez un notebook dans le pipeline qui appelle une fonction exit(), l’activité du notebook retourne une valeur de sortie, termine l’exécution du pipeline et arrête la session Spark. Ne limitez pas la fonction exit() autour d’un try/catch car cette exception NotebookExit doit se propager pour que le pipeline obtienne la valeur de retour.

  • Lorsque vous appelez une fonction exit() dans un carnet référencée, Fabric Spark arrête l’exécution ultérieure du carnet référencé et continue d’exécuter les cellules suivantes du carnet principal qui appelle la fonction run(). Par exemple : Notebook1 possède trois cellules et appelle une fonction exit() dans la deuxième cellule. Notebook2 possède cinq cellules et appelle run(notebook1) dans la troisième cellule. Lorsque vous exécutez Notebook2, Notebook1 s’arrête à la deuxième cellule lorsque vous atteignez la fonction exit(). Notebook2 continue à exécuter ses quatrième et cinquième cellules.

mssparkutils.notebook.exit("value string")

Par exemple :

Le notebook Sample1 avec les deux cellules suivantes :

  • La cellule 1 définit un paramètre d’entrée dont la valeur par défaut est définie sur 10.

  • La cellule 2 quitte le notebook avec l'entrée comme valeur de sortie.

Capture d’écran d’un exemple de notebook de la fonction exit.

Vous pouvez exécuter Sample1 dans un autre bloc-notes avec les valeurs par défaut :

exitVal = mssparkutils.notebook.run("Sample1")
print (exitVal)

Sortie:

Notebook executed successfully with exit value 10

Vous pouvez exécuter Sample1 dans un autre notebook et définir la valeur d’entrée sur 20 :

exitVal = mssparkutils.notebook.run("Sample1", 90, {"input": 20 })
print (exitVal)

Sortie:

Notebook executed successfully with exit value 20

Utilitaires d’informations d’identification

Vous pouvez utiliser les utilitaires d’identifiants MSSparkUtils pour obtenir des tokens d’accès et gérer des secrets dans Azure Key Vault.

Exécutez la commande suivante pour obtenir une vue d’ensemble des méthodes disponibles :

mssparkutils.credentials.help()

Sortie:

getToken(audience, name): returns AAD token for a given audience, name (optional)
getSecret(keyvault_endpoint, secret_name): returns secret for a given Key Vault and secret name

Obtenir un jeton

getTokenrenvoie un jeton Microsoft Entra pour un public donné et un nom (optionnel). La liste suivante présente les touches d’audience actuellement disponibles :

  • Ressource pour le public de stockage : storage
  • Ressource Power BI :pbi
  • Azure Key Vault Resource :keyvault
  • Ressource de base de données KQL Synapse RTA : kusto

Exécutez la commande suivante pour obtenir le jeton :

mssparkutils.credentials.getToken('audience Key')

Obtenir le secret à l’aide des identifiants de l’utilisateur

getSecret renvoie un secret Azure Key Vault pour un point de terminaison Azure Key Vault et un nom de secret donnés, à l’aide des informations d’identification de l’utilisateur.

mssparkutils.credentials.getSecret('https://<name>.vault.azure.net/', 'secret name')

Montage et démontage de fichiers

Fabric prend en charge les scénarios de montage suivants dans le package Utilitaires Microsoft Spark. Vous pouvez utiliser les API mount, unmount, getMountPath() et mounts() pour connecter le stockage distant (Azure Data Lake Storage Gen2) à tous les nœuds en activité (nœuds pilotes et nœuds ouvriers). Une fois le point de montage de stockage en place, utilisez l’API de fichier local pour accéder aux données comme si elles étaient stockées dans le système de fichiers local.

Comment monter un compte Azure Data Lake Storage Gen2

L’exemple suivant montre comment monter Azure Data Lake Storage Gen2. Le montage de Stockage Blob fonctionne de la même façon.

Cet exemple suppose que vous disposez d’un compte Data Lake Storage Gen2 nommé storegen2et que le compte possède un conteneur nommé mycontainer que vous souhaitez monter sur /test dans votre session Spark de notebook.

Capture d’écran montrant où sélectionner un conteneur à monter.

Pour monter le conteneur nommé mycontainer, mssparkutils vérifie d’abord si vous avez la permission d’accéder au conteneur. Fabric prend en compte trois méthodes d’authentification pour l’opération de montage de déclenchement : le jeton Microsoft Entra (par défaut et recommandé), accountKey et sastoken. Pour plus d’informations sur l’authentification par jeton Microsoft Entra et l’API actuellenotebookutils, consultez Montage et démontage de fichiers avec NotebookUtils pour Fabric.

Monter en utilisant un jeton de signature d’accès partagé ou une clé de compte

MSSparkUtils permet de transmettre explicitement une clé de compte ou un jeton de signature d’accès partagé (SAP) en tant que paramètre pour monter la cible.

Pour des raisons de sécurité, nous vous recommandons de stocker des clés de compte ou des jetons SAP dans Azure Key Vault (comme l’illustre la capture d’écran suivant). Vous pouvez ensuite les récupérer à l’aide de l’API mssparkutils.credentials.getSecret. Pour plus d’informations sur Azure Key Vault, consultez À propos des clés de compte de stockage managées Azure Key Vault.

Capture d'écran indiquant les emplacements de stockage des secrets dans un coffre Azure Key Vault.

Exemple de code pour la méthode accountKey :

from notebookutils import mssparkutils  
# get access token for keyvault resource
# you can also use full audience here like https://vault.azure.net
accountKey = mssparkutils.credentials.getSecret("<vaultURI>", "<secretName>")
mssparkutils.fs.mount(  
    "abfss://mycontainer@<accountname>.dfs.core.windows.net",  
    "/test",  
    {"accountKey":accountKey}
)

Exemple de code pour sastoken :

from notebookutils import mssparkutils  
# get access token for keyvault resource
# you can also use full audience here like https://vault.azure.net
sasToken = mssparkutils.credentials.getSecret("<vaultURI>", "<secretName>")
mssparkutils.fs.mount(  
    "abfss://mycontainer@<accountname>.dfs.core.windows.net",  
    "/test",  
    {"sasToken":sasToken}
)

Remarque

Il se peut que vous deviez importer mssparkutils s’il n’est pas disponible :

from notebookutils import mssparkutils

Paramètres de montage :

  • fileCacheTimeout: Les blobs sont mis en cache dans le dossier temporaire local pendant 120 secondes par défaut. Pendant ce temps, blobfuse ne vérifie pas si le fichier est à jour. Définissez ce paramètre pour changer le délai d’attente par défaut. Lorsque plusieurs clients modifient des fichiers en même temps, pour éviter les incohérences entre les fichiers locaux et distants, nous recommandons de raccourcir le temps de cache, voire de le réduire à zéro, et de toujours récupérer les fichiers les plus récents depuis le serveur.
  • timeout: Le délai d’attente pour l’opération de monture est de 120 secondes par défaut. Définissez ce paramètre pour changer le délai d’attente par défaut. Lorsqu’il y a trop d’exécuteurs ou que le délai d’expiration du montage est atteint, nous recommandons d’augmenter cette valeur.

Vous pouvez utiliser ces paramètres comme suit :

mssparkutils.fs.mount(
   "abfss://mycontainer@<accountname>.dfs.core.windows.net",
   "/test",
   {"fileCacheTimeout": 120, "timeout": 120}
)

Remarque

Par souci de sécurité, ne stockez pas d’informations d’identification dans du code. Pour mieux protéger vos identifiants, le secret est masqué dans la sortie du notebook. Pour plus d’informations, consultez Suppression des secrets.

Guide pratique pour monter un lakehouse

Exemple de code pour monter une maison lacustre vers /test:

from notebookutils import mssparkutils 
mssparkutils.fs.mount( 
 "abfss://<workspace_id>@onelake.dfs.fabric.microsoft.com/<lakehouse_id>", 
 "/test"
)

Remarque

Le montage d’un point de terminaison régional n’est pas pris en charge. Fabric prend uniquement en charge le montage du point de terminaison global, onelake.dfs.fabric.microsoft.com.

Accédez aux fichiers sous le point de montage en utilisant l’API fs mssparkutils

Le principal objectif de l’opération de montage est de vous permettre d’accéder aux données stockées dans un compte de stockage distant en utilisant une API locale du système de fichiers. Vous pouvez également accéder aux données à l’aide de l’API mssparkutils fs avec un chemin d’accès monté en tant que paramètre. Le format de ce chemin est un peu différent.

Supposons que vous ayez monté le conteneur mycontainer de /test Data Lake Storage Gen2 en utilisant l’API mount. Lorsque vous accédez aux données en utilisant une API locale du système de fichiers, le format de chemin est le suivant :

/synfs/notebook/{sessionId}/test/{filename}

Lorsque vous souhaitez accéder aux données en utilisant l’API mssparkutils fs , nous vous recommandons d’utiliser getMountPath() pour obtenir le chemin précis :

path = mssparkutils.fs.getMountPath("/test")
  • Lister des répertoires :

    mssparkutils.fs.ls(f"file://{mssparkutils.fs.getMountPath('/test')}")
    
  • Lisez le contenu du fichier :

    mssparkutils.fs.head(f"file://{mssparkutils.fs.getMountPath('/test')}/myFile.txt")
    
  • Créer un répertoire :

    mssparkutils.fs.mkdirs(f"file://{mssparkutils.fs.getMountPath('/test')}/newdir")
    

Accéder aux fichiers sous le point de montage via le chemin d’accès local

Vous pouvez facilement lire et écrire les fichiers au point de montage à l’aide du système de fichiers standard. Voici un exemple Python :

#File read
with open(mssparkutils.fs.getMountPath('/test2') + "/myFile.txt", "r") as f:
    print(f.read())
#File write
with open(mssparkutils.fs.getMountPath('/test2') + "/myFile.txt", "w") as f:
    print(f.write("dummy data"))

Guide pratique pour vérifier les points de montage existants

Vous pouvez utiliser l’API mssparkutils.fs.mounts() pour vérifier toutes les informations de point de montage existantes :

mssparkutils.fs.mounts()

Comment démonter le point de montage

Utilisez le code suivant pour démonter votre point de montage (/test dans cet exemple) :

mssparkutils.fs.unmount("/test")

Limitations connues

  • Le montage actuel est une configuration au niveau de la tâche. Nous vous recommandons d’utiliser l’API mounts pour vérifier si un point de montage existe ou n’est pas disponible.

  • Le mécanisme de démontage n’est pas automatique. Une fois l’exécution de l’application terminée, pour démonter le point de montage et libérer l’espace disque, vous devez appeler explicitement une API de démontage dans votre code. Sinon, le point de montage existe toujours dans le nœud une fois l’exécution de l’application terminée.

  • Le montage d’un compte de stockage Azure Data Lake Storage Gen1 n’est pas pris en charge.

Utilitaires de lakehouse

Le mssparkutils.lakehouse module fournit des utilités pour gérer les éléments de la maison du lac. Ces utilitaires facilitent la création, la récupération, la mise à jour et la suppression d’éléments de lakehouse.

Remarque

Les API Lakehouse sont prises en charge uniquement sur la version 1.2 ou ultérieure de Runtime.

Vue d’ensemble des méthodes

Les méthodes suivantes sont disponibles dans le mssparkutils.lakehouse module :

# Create a new Lakehouse artifact
create(name: String, description: String = "", workspaceId: String = ""): Artifact

# Retrieve a Lakehouse artifact
get(name: String, workspaceId: String = ""): Artifact

# Update an existing Lakehouse artifact
update(name: String, newName: String, description: String = "", workspaceId: String = ""): Artifact

# Delete a Lakehouse artifact
delete(name: String, workspaceId: String = ""): Boolean

# List all Lakehouse artifacts
list(workspaceId: String = ""): Array[Artifact]

Exemples d'utilisation

Pour utiliser ces méthodes efficacement, considérez les exemples d’utilisation suivants :

Créer un objet de maison lacustre

artifact = mssparkutils.lakehouse.create("artifact_name", "Description of the artifact", "optional_workspace_id")

Récupérer un objet de la maison du lac

artifact = mssparkutils.lakehouse.get("artifact_name", "optional_workspace_id")

Mise à jour d’un élément de la maison du lac

updated_artifact = mssparkutils.lakehouse.update("old_name", "new_name", "Updated description", "optional_workspace_id")

Suppression d’un objet de la maison du lac

is_deleted = mssparkutils.lakehouse.delete("artifact_name", "optional_workspace_id")

Liste des éléments de la maison du lac

artifacts_list = mssparkutils.lakehouse.list("optional_workspace_id")

Informations supplémentaires

Pour des informations plus détaillées sur chaque méthode et ses paramètres, utilisez la mssparkutils.lakehouse.help("methodName") fonction.

En utilisant les utilitaires Lakehouse de MSSparkUtils, vous pouvez gérer plus efficacement vos éléments Lakehouse et intégrer cette gestion dans vos pipelines Fabric, améliorant ainsi votre expérience globale de gestion des données.

Explorez ces utilités et intégrez-les dans vos flux de travail Fabric pour une gestion fluide des articles de la maison du lac.

Utilitaires de runtime

Afficher les informations de contexte de la session

En utilisant mssparkutils.runtime.context, vous pouvez obtenir les informations de contexte de la session active actuelle, notamment le nom du notebook, le lakehouse par défaut, les informations sur l’espace de travail, s’il s’agit d’une exécution de pipeline, entre autres.

mssparkutils.runtime.context

Remarque

mssparkutils.env n'est pas officiellement pris en charge sur Fabric. Utilisez notebookutils.runtime.context comme alternative.

Problème connu

Lorsque vous utilisez une version d’exécution ultérieure à la 1.2 et que vous exécutezmssparkutils.help(), les APIfabricClient, Warehouse et workspace listées ne sont pas actuellement prises en charge.