Aller au contenu principal

Trigger jobs à l'arrivée de nouveaux fichiers

Vous pouvez utiliser des triggers d'arrivée de fichiers pour déclencher l'exécution de votre job lorsque de nouveaux fichiers arrivent dans un emplacement externe tel qu'Amazon S3, le stockage Azure ou Google Cloud Storage. Cette fonctionnalité est utile lorsque l'efficacité d'un job planifié est compromise par des arrivées de nouvelles données irrégulières.

Comment fonctionnent les Trigger d'arrivée de fichier

Les Trigger d'arrivée de fichiers s'efforcent de vérifier la présence de nouveaux fichiers chaque minute, bien que cela puisse être affecté par les performances du stockage cloud sous-jacent. Les Trigger d'arrivée de fichiers n'entraînent pas de coûts supplémentaires autres que les coûts du fournisseur de cloud associés à la liste des fichiers dans l'emplacement de stockage.

Un trigger d'arrivée de fichiers peut être configuré pour surveiller la racine d'un emplacement externe ou d'un volume Unity Catalog, ou un sous-chemin d'un emplacement externe ou d'un volume. Par exemple, pour le volume Unity Catalog /Volumes/mycatalog/myschema/myvolume/, les chemins suivants sont valides pour un Trigger d'arrivée de fichier :

/Volumes/mycatalog/myschema/myvolume/
/Volumes/mycatalog/myschema/myvolume/mydirectory/

Un Trigger d'arrivée de fichier vérifie de manière récursive les nouveaux fichiers dans tous les sous-répertoires de l'emplacement configuré. Par exemple, vous créez un Trigger d'arrivée de fichier pour l'emplacement /Volumes/mycatalog/myschema/myvolume/mydirectory/, et cet emplacement contient les sous-répertoires suivants :

/Volumes/mycatalog/myschema/myvolume/mydirectory/subdirA
/Volumes/mycatalog/myschema/myvolume/mydirectory/subdirB
/Volumes/mycatalog/myschema/myvolume/mydirectory/subdirC/subdirD

Le trigger vérifie les nouveaux fichiers dans mydirectory, subdirA, subdirB, subdirC et subdirC/subdirD.

File arrival Trigger with file events

Pour des performances optimales, l'emplacement externe doit être activé pour les événements de fichier. Lorsque les événements de fichier sont activés pour un emplacement externe, Databricks utilise un service interne pour suivre les métadonnées d'ingestion en traitant les notifications de modification des fournisseurs de cloud. Ce service conserve les métadonnées des fichiers les plus récents créés ou mis à jour pendant une période de rétention glissante déterminée par le service, ce qui améliore l'efficacité du traitement des fichiers.

Quelques minutes après l'activation des événements de fichiers sur un emplacement externe, les triggers d'arrivée de fichiers existants qui surveillent les chemins couverts par cet emplacement externe start à bénéficier de l'activation des événements de fichiers, et les nouveaux triggers en bénéficient en quelques secondes.

Pour plus d'informations concernant les avantages en termes de performances et de capacité des événements de fichiers sur les emplacements externes, consultez Limitations. Pour les questions fréquentes sur les événements de fichiers, consultez la FAQ sur les événements de fichiers.

Avant de commencer

Les éléments suivants sont requis pour utiliser les Triggers d'arrivée de fichiers :

Ajouter un trigger d'arrivée de fichier

Pour ajouter un Trigger d'arrivée de fichier à un Job :

  1. Dans la barre latérale de votre workspace Databricks, cliquez sur Tâches & Pipelines .
  2. Facultativement, sélectionnez les filtres **Jobs** et **Appartenant à moi**.
  3. Cliquez sur le **Link** **Nom** de votre Job.
  4. Dans le volet Détails du Job à droite, cliquez sur Ajouter un trigger .
  5. Dans Type de Trigger , sélectionnez Arrivée de fichier .
  6. Dans Emplacement de stockage , saisissez l'URL de la racine ou d'un sous-chemin d'un emplacement externe de Unity Catalog ou de la racine ou d'un sous-chemin d'un volume Unity Catalog à surveiller.
  7. (Facultatif) Configurez les options avancées ( Temps minimal entre les Trigger en secondes et Délai d'attente après la dernière modification en secondes ) pour contrôler la fréquence de Trigger des exécutions. Pour des exemples de configuration, consultez Contrôler la fréquence à laquelle les exécutions sont Trigger.
  8. Pour valider la configuration, cliquez sur **Tester la connexion**.
  9. Cliquez sur Enregistrer .

Pour modifier, suspendre ou supprimer ce Trigger ultérieurement, utilisez la section Plannings & Triggers du volet Détails du Job . Consultez Gérer un Trigger existant.

Contrôler la fréquence à laquelle les Trigger sont déclenchées

Deux options avancées sur un Trigger d'arrivée de fichier contrôlent la manière dont les arrivées de fichiers se traduisent en exécutions de Job. Ces options appliquent deux modèles courants de contrôle de taux : un refroidissement et une élimination des rebonds .

  • **Délai minimal entre les Trigger en secondes** : Limite le Job à une exécution maximale par cet intervalle (une période de latence entre les exécutions). Une fois qu'une exécution est terminée, les fichiers qui arrivent pendant la période de latence ne start pas une nouvelle exécution tant que l'intervalle n'est pas écoulé. Utilisez cette option pour plafonner la fréquence de création des exécutions afin que les arrivées fréquentes ne créent pas d'exécutions consécutives.
  • Temps d'attente après le dernier changement en secondes : attend ce délai après l'arrivée du fichier le plus récent avant de démarrer une exécution, et chaque nouvelle arrivée réinitialise le minuteur (antibond). Utilisez cette option lorsque les fichiers arrivent par batchs et que vous souhaitez traiter l'intégralité du batch en une seule exécution après que tous les fichiers aient été déposés.

Vous pouvez définir l’une ou l’autre option seule, ou les deux ensemble. Voir les exemples suivants.

Exécuter au plus toutes les 15 minutes

Pour créer des exécutions à mesure que les fichiers arrivent, mais pas plus fréquemment que toutes les 15 minutes, définissez l'option avancée suivante :

  • Délai minimal entre les déclencheurs en secondes : 900

Une fois chaque exécution terminée, le Trigger attend 900 secondes (15 minutes) avant de start une autre exécution, même si des fichiers continuent d’arriver. Ceci limite la création d'exécutions à une seule exécution toutes les 15 minutes au maximum.

Attendez qu'un batch complet arrive

Lorsque les fichiers arrivent par lots et que vous souhaitez traiter chaque lot en une seule exécution, définissez Attendre après la dernière modification en secondes sur une valeur inférieure à l'intervalle entre les lots, mais supérieure à l'intervalle entre les fichiers au sein d'un lot. Par exemple, si un nouveau batch start toutes les 5 minutes environ, définissez l'option avancée suivante :

  • Attendre après la dernière modification (en secondes) : 60

Chaque nouveau fichier Reset le minuteur, ainsi le Trigger start une exécution seulement après 60 secondes sans nouvelle arrivée. Cette configuration suppose que les fichiers au sein d'un batch arrivent à 60 secondes d'intervalle les uns des autres, afin que le minuteur n'expire pas au milieu d'un batch, et que les batchs sont espacés de plus de 60 secondes, afin que les batchs consécutifs ne Merge pas en une seule exécution.

Limiter la fréquence et attendre les batchs complets

Vous pouvez combiner les deux options lorsque vous voulez plafonner la fréquence de création des exécutions et éviter de démarrer une exécution au milieu d'un batch. Par exemple :

  • Délai minimal entre les déclencheurs en secondes : 900
  • Attendre après la dernière modification (en secondes) : 60

Avec cette configuration, le Trigger attend la fin du débarquement d'un batch (60 secondes sans nouveaux fichiers) avant de start une exécution, et il ne start pas plus d'une exécution toutes les 15 minutes.

Découvrir et traiter les fichiers à l'arrivée

Pour traiter les fichiers qui ont déclenché les Trigger d'arrivée de fichier, vous pouvez utiliser Auto Loader. Auto Loader traite les nouveaux fichiers de manière incrémentielle et efficace avec des garanties d’exactement une fois. Par exemple, utilisez l’extrait ci-dessous pour charger des fichiers dans une table Delta.

Pour utiliser cette solution, créez un Job avec un Trigger d’arrivée de fichier et ajoutez un Notebook contenant le code ci-dessous. Remplacez chaque espace réservé [REPLACE] par la valeur appropriée.

Python
# Configuration
file_location = "[REPLACE]" # The same URL configured for the file arrival trigger.
checkpoint_location = "[REPLACE]" # a separate URL (outside `file_location`) used to store the Auto Loader checkpoint, which enables exactly-once processing.
sink_table = "[REPLACE]" # Delta table to write to

# Use Auto Loader to discover new files.
# Do not modify code below this line
streamingQuery = spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "json") \
.option("cloudFiles.schemaLocation", checkpoint_location) \
.option("cloudFiles.useManagedFileEvents","true") \
.load(file_location) \
.writeStream \
.option("checkpointLocation", checkpoint_location) \
.trigger(availableNow = True) \
.toTable(sink_table)

Si vous devez traiter de nouveaux fichiers avec une logique personnalisée et que vous souhaitez uniquement découvrir l'URL des nouveaux fichiers, vous pouvez utiliser foreachBatch à la place, comme indiqué dans l'extrait de code ci-dessous. Notez que foreachBatch fournit uniquement des garanties de traitement au moins une fois. Pour plus d'information sur l'utilisation de foreachBatch, consultez Utiliser foreachBatch pour écrire dans des récepteurs de données arbitraires

Python
# Configuration
file_location = "[REPLACE]" # The same URL configured for the file arrival trigger.
checkpoint_location = "[REPLACE]" # a separate URL (outside `file_location`) used to store the Auto Loader checkpoint, which enables exactly-once processing.

def process_batch(batch_df, batch_id):
file_url = batch_df.select("path").collect()[0].path
# [REPLACE] Your custom function for processing newly arrived files


# Use Auto Loader to discover new files.
# Do not modify code below this line
streamingQuery = spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "binaryFile") \
.option("cloudFiles.useManagedFileEvents","true") \
.load(file_location) \
.drop("content") \
.writeStream \
.foreachBatch(process_batch) \
.option("checkpointLocation", checkpoint_location) \
.trigger(availableNow = True) \
.start()

Recevoir des notifications de Trigger d'arrivée de fichiers ayant échoué

Pour être averti si un Trigger d’arrivée de fichiers ne parvient pas à être évalué, configurez les notifications par e-mail ou par destination système en cas d’échec du Job. Voir Ajouter des notifications à un job.

Limitations

  • Seuls les nouveaux fichiers déclenchent des exécutions. L'écrasement d'un fichier existant par un fichier du même nom ne Trigger pas d'exécution.

  • Le chemin utilisé pour un Trigger d'arrivée de fichier ne doit pas contenir de tables externes ou d'emplacements gérés de catalogues et de schémas.

  • Le chemin utilisé pour un trigger d'arrivée de fichier ne peut pas contenir de caractères génériques, par exemple, * ou ?.

  • Si l'emplacement de stockage est configuré comme un emplacement externe dans Unity Catalog et que cet emplacement externe est activé pour les événements de fichier:

    • Il n'y a pas de limites sur le nombre de fichiers dans l'emplacement de stockage.

    • Les Trigger peuvent générer une erreur en cas de délai d'expiration lorsqu'il y a trop de mises à jour de fichiers superflues.

      Lorsqu'un trigger d'arrivée de fichiers est défini sur un sous-chemin d'un emplacement externe ou d'un volume Unity Catalog, les changements en dehors de ce sous-chemin, tels qu'à la racine de l'emplacement externe, peuvent augmenter la quantité de métadonnées que le trigger doit traiter. Dans les environnements à forte évolution, cela peut entraîner le dépassement de la limite de temps de traitement du Trigger et un état d'erreur.

      Pour éviter cela, créez un volume Unity Catalog qui correspond spécifiquement au sous-répertoire que vous souhaitez surveiller et définissez le Trigger d’arrivée de fichiers à la racine de ce volume. Cette approche isole votre chemin cible comme racine effective du Trigger, réduisant les changements de niveau racine non liés et empêchant le Trigger de passer à un état d'erreur.

    • Si un fichier existant est modifié et que ses métadonnées tombent en dehors de la période de rétention glissante, cette modification est traitée comme une nouvelle arrivée de fichier, déclenchant une exécution de Job. Vous pouvez éviter cela en ingérant uniquement des fichiers immuables, ou vous pouvez utiliser des déclencheurs d'arrivée de fichier avec Auto Loader pour suivre la progression de l'ingestion.

  • Si l'emplacement de stockage n'est pas activé pour les événements de fichier :

    • Un maximum de 50 jobs peut être configuré avec un trigger d'arrivée de fichier sur de tels emplacements dans un workspace Databricks.
    • L'emplacement de stockage peut contenir jusqu'à 10 000 fichiers. Si l'emplacement de stockage configuré est un sous-chemin d'un emplacement externe ou d'un volume Unity Catalog, la limite de 10 000 fichiers s'applique au sous-chemin et non à la racine de l'emplacement de stockage. Par exemple, la racine de l'emplacement de stockage peut contenir plus de 10 000 fichiers dans ses sous-répertoires, mais le sous-répertoire configuré ne doit pas dépasser la limite de 10 000 fichiers.

Voir aussi limitations des événements de fichier.

Trigger d'arrivée de fichiers sur des chemins non existants dans les emplacements externes S3 et GCS

Lorsque le répertoire configuré n'existe pas ou est supprimé d'Amazon S3 ou de Google Cloud Storage, les triggers d'arrivée de fichiers continuent d'être évalués sans erreur. Ce comportement se produit parce que S3 et GCS ne font pas de distinction entre les répertoires non existants, supprimés et vides.

Par conséquent, un trigger d'arrivée de fichier monitoring un chemin de répertoire non existant ou supprimé n'échoue pas et ne génère pas de notification d'erreur. Le Trigger continue d'évaluer, ne trouve aucun fichier et ne déclenche aucune exécution de Job tant que des fichiers ne sont pas ajoutés à ce chemin à nouveau. Il s'agit d'un comportement attendu et non d'une condition d'erreur.