Aller au contenu principal

Lectures et écritures de streaming de tables Delta Lake

Cette page décrit comment utiliser les tables Delta Lake comme sources et puits de données pour Spark Structured Streaming avec readStream et writeStream. Delta Lake résout les problèmes courants de performance et de fiabilité des systèmes de streaming et des fichiers. Les avantages comprennent :

  • Regroupez les petits fichiers produits par l'ingestion à faible latence et améliorez les performances.
  • Maintenir le traitement « exactly-once » avec plus d’un Stream (ou des batch Jobs concurrents).
  • Détecter efficacement les nouveaux fichiers lorsque des fichiers sont utilisés comme source de Stream.

Pour savoir comment charger des données à l’aide de tables de streaming dans Databricks SQL, consultez Utiliser les tables de streaming autonomes.

Pour les jointures Stream-statiques avec Delta Lake, consultez les jointures Stream-statiques.

Pour obtenir une liste complète des options DataStreamReader et DataStreamWriter pour Delta Lake, consultez DataStreamReader options Delta Lake et DataStreamWriter options Delta Lake.

attention

Si vous utilisez une table Delta Lake comme source de streaming, la query de streaming doit s'exécuter au moins une fois dans la fenêtre de rétention de la table source. Les fenêtres de rétention default sont de 7 jours pour les fichiers de données supprimés VACUUMet de 30 jours pour les Logs de transaction (logRetentionDuration). Si une requête prend du retard sur ces fenêtres, elle échoue avec DELTA_FILE_NOT_FOUND_DETAILED et doit être reset avec un refresh complet.

Ne *définissez pas* spark.sql.files.ignoreMissingFiles sur true comme solution de contournement, car cette configuration produit silencieusement des résultats incorrects. Si la planification d'un stream ne peut pas suivre les fenêtres de rétention par default, augmentez plutôt la rétention de la table source.

Utiliser les tables Delta Lake comme récepteur

Vous pouvez écrire des données dans une table Delta Lake en utilisant Structured Streaming. Le Log des transactions Delta Lake garantit un traitement exactement une fois, même lorsque d'autres Streams ou batch requêtes sont exécutés simultanément sur la table.

Lorsque vous écrivez dans une table Delta Lake à l’aide d’un récepteur Structured Streaming, vous pouvez voir des commits vides avec epochId = -1. Ceux-ci sont attendus et se produisent généralement :

  • Lors du premier batch de chaque exécution de la query de streaming (cela se produit à chaque batch pour Trigger.AvailableNow).
  • Lorsqu'un schéma est modifié (tel que l'ajout d'une colonne).

Ces commits vides sont intentionnels et n'indiquent pas d'erreur. Elles n'affectent pas l'exactitude ou les performances de la query de manière significative.

remarque

La fonction VACUUM de Delta Lake supprime tous les fichiers non gérés par Delta Lake, mais ignore les répertoires commençant par _. Vous pouvez stocker en toute sécurité les points de contrôle aux côtés d'autres données et métadonnées pour une table Delta Lake en utilisant une structure de répertoires telle que <table-name>/_checkpoints.

Surveiller le backlog avec des métriques

Utilisez les métriques suivantes pour surveiller l'arriéré d'un processus de query en streaming:

  • numBytesOutstanding: Nombre d'octets restant à traiter dans le backlog.
  • numFilesOutstanding– Nombre de fichiers à traiter dans le backlog.
  • numNewListedFiles: Nombre de fichiers Delta Lake listés pour calculer l'arriéré pour ce batch.
  • backlogEndOffset: La version de la table Delta Lake utilisée pour calculer le backlog.

Dans un Notebook, consultez ces métriques sous l'onglet données brutes du tableau de bord de progression de la query de streaming :

JSON
{
"sources": [
{
"description": "DeltaSource[file:/path/to/source]",
"metrics": {
"numBytesOutstanding": "3456",
"numFilesOutstanding": "8"
}
}
]
}

Mode d’ajout

By default, streams run in append mode and only add new records to the table.

Utilisez la méthode toTable lors du streaming vers des tables :

Python
(events.writeStream
.outputMode("append")
.option("checkpointLocation", "/tmp/delta/events/_checkpoints/")
.toTable("events")
)

Mode Complet

Utilisez Structured Streaming avec le mode complet pour remplacer l'intégralité de la table après chaque batch. Par exemple, vous pouvez continuellement mettre à jour une table récapitulative agrégée des événements par client :

Python
(spark.readStream
.table("events")
.groupBy("customerId")
.count()
.writeStream
.outputMode("complete")
.option("checkpointLocation", "/tmp/delta/eventsByCustomer/_checkpoints/")
.toTable("events_by_customer")
)

Pour les applications sans exigences strictes en matière de latence, vous pouvez économiser des ressources de calcul et des coûts avec des Trigger uniques tels que AvailableNow. Par exemple, utilisez ce trigger pour mettre à jour les tables d'agrégation récapitulatives selon un calendrier donné, en traitant uniquement les nouvelles données arrivées depuis la dernière mise à jour. Consultez AvailableNow: Traitement par batch incrémentiel.

Gérer les modifications apportées aux tables Delta Lake sources

Structured Streaming lit de manière incrémentielle les tables Delta Lake. Lorsqu'une query de streaming lit à partir d'une table Delta Lake, les nouveaux enregistrements sont traités de manière idempotente à mesure que de nouvelles versions de table sont commit dans la table source. Structured Streaming n'accepte que les entrées d'ajout et lève une exception si des modifications se produisent sur la table source Delta Lake. Par exemple, si une opérations UPDATE, DELETE, MERGE INTO ou OVERWRITE modifie une table source Delta Lake qui est lue par une query de streaming, le Stream échoue avec une erreur.

Il existe quatre approches typiques pour gérer les modifications en amont des tables source Delta Lake, en fonction de votre cas d'utilisation. Vous trouverez ci-dessous un tableau de référence et des détails sur chacun d’entre eux :

Approche

Avantages

Consommation

skipChangeCommits

Simple, ne vous demande pas d'écrire une logique complexe. Utile pour le traitement en mode ajout uniquement où les modifications en amont sont gérées séparément, ou pour gérer temporairement un mauvais enregistrement.

Ne propage pas les changements et ne traite que les ajouts.

refresh complète

Également simple, cela ne nécessite pas d'écrire une logique complexe. Utile pour les petits datasets avec des modifications amont rares.

Coûteux pour les grands datasets. Nécessite le retraitement de toutes les tables en aval.

Flux de données de modification

Traiter tous les types de modification (insertions, mises à jour et suppressions). Databricks recommande le streaming à partir du flux CDC d'une table Delta Lake plutôt que directement à partir de la table, dans la mesure du possible.

Nécessite d'écrire une logique plus complexe pour gérer chaque type de modification.

Vues matérialisées

Alternative simple au Structured Streaming qui offre une propagation automatique des modifications.

Latence plus élevée. Disponible uniquement dans LakeFlow Pipelines et Databricks SQL.

Approche

Avantages

Consommation

skipChangeCommits

Simple, ne vous demande pas d'écrire une logique complexe. Utile pour le traitement en mode ajout uniquement où les modifications en amont sont gérées séparément, ou pour gérer temporairement un mauvais enregistrement.

Ne propage pas les changements et ne traite que les ajouts.

refresh complète

Également simple, cela ne nécessite pas d'écrire une logique complexe. Utile pour les petits datasets avec des modifications amont rares.

Coûteux pour les grands datasets. Nécessite le retraitement de toutes les tables en aval.

Flux de données de modification

Traiter tous les types de modification (insertions, mises à jour et suppressions). Databricks recommande le streaming à partir du flux CDC d'une table Delta Lake plutôt que directement à partir de la table, dans la mesure du possible.

Nécessite d'écrire une logique plus complexe pour gérer chaque type de modification.

Vues matérialisées

Alternative simple au Structured Streaming qui offre une propagation automatique des modifications.

Latence plus élevée. Disponible uniquement dans LakeFlow Pipelines et Databricks SQL.

Ignorer les commits de modification en amont avec skipChangeCommits

Définissez skipChangeCommits pour ignorer les transactions qui suppriment ou modifient des enregistrements existants, et pour traiter uniquement les ajouts. Ceci est utile lorsque les modifications aux données existantes n'ont pas besoin d'être propagées via le Stream, ou lorsque vous préférez une logique distincte pour gérer ces modifications. Vous pouvez activer et désactiver skipChangeCommits si vous devez ignorer temporairement les modifications ponctuelles.

Databricks recommande d’utiliser skipChangeCommits pour la plupart des workloads qui n’utilisent pas de flux de données modifiés.

Python
(spark.readStream
.option("skipChangeCommits", "true")
.table("source_table")
)
important

Si le schéma d'une table Delta Lake change après le début d'une lecture en streaming sur la table, la query échoue. Pour la plupart des modifications de schéma, vous pouvez redémarrer le Stream afin de résoudre les incohérences de schéma et de poursuivre le traitement.

Dans Databricks Runtime 12.2 LTS et versions antérieures, vous ne pouvez pas effectuer de Stream à partir d'une table Delta Lake avec le mappage de colonnes activé qui a subi une évolution des schémas non additive, telle que le renommage ou la suppression de colonnes. Pour plus de détails, consultez Mappage des colonnes et streaming.

remarque

Dans Databricks Runtime 12.2 LTS et versions ultérieures, skipChangeCommits remplace ignoreChanges. Dans Databricks Runtime 11.3 LTS et versions antérieures, ignoreChanges est la seule option prise en charge. Consultez l’option héritée : ignoreChanges pour plus de détails.

Option héritée : ignoreDeletes

ignoreDeletes est une option héritée qui ne gère que les transactions qui suppriment des données aux limites de partition (c'est-à-dire, les suppressions complètes de partition). Si vous devez gérer des suppressions hors partition, des mises à jour ou d'autres modifications, utilisez skipChangeCommits à la place.

Python
(spark.readStream
.option("ignoreDeletes", "true")
.table("user_events")
)

Option héritée : ignoreChanges

ignoreChanges est disponible dans Databricks Runtime 11.3 LTS et les versions inférieures. Dans Databricks Runtime 12.2 LTS et les versions supérieures, il est remplacé par skipChangeCommits.

Avec ignoreChanges activé, les fichiers de données réécrits dans la table source sont réémis après une opération de modification de données telle que UPDATE, MERGE INTO, DELETE (au sein des partitions) ou OVERWRITE. Les lignes inchangées sont souvent émises avec les nouvelles lignes, les consommateurs en aval doivent donc être capables de gérer les doublons. Les suppressions ne sont pas propagées en aval. ignoreChanges prime sur ignoreDeletes.

En revanche, skipChangeCommits ignore entièrement les Opérations de modification de fichiers. Les fichiers de données réécrits dans la table source en raison d'Opérations de modification de données telles que UPDATE, MERGE INTO, DELETE et OVERWRITE sont entièrement ignorés. Pour refléter les modifications dans les tables source de stream, vous devez implémenter une logique distincte pour propager ces modifications.

Databricks recommande d'utiliser skipChangeCommits pour toutes les nouvelles workloads. Pour migrer une charge de travail de ignoreChanges vers skipChangeCommits, refactorisez votre logique de streaming.

Full refresh des tables en aval

Si les changements en amont sont rares et que les données sont suffisamment petites pour être retraitées, vous pouvez supprimer le point de contrôle de streaming et la table de sortie, puis redémarrer le stream depuis le début. Cela fait que le Stream retraitera toutes les données de la table source. Notez que cette approche nécessite également de retraiter toutes les tables en aval qui dépendent de la sortie de ce stream.

Cette approche convient mieux aux petits datasets ou aux workloads où les modifications en amont sont peu fréquentes et où le coût d'une full refresh est acceptable.

Utiliser le flux de données de modification

Pour les charges de travail qui traitent tous les types de modifications (insertions, mises à jour et suppressions), utilisez le flux de données de modification Delta Lake. Le flux de données de modification enregistre les modifications au niveau des lignes dans une table Delta Lake, ce qui vous permet de Stream ces modifications et d'écrire une logique pour gérer chaque type de modification dans les tables en aval. C'est l'approche la plus robuste, car votre code gère explicitement chaque type d'événement de modification. Consultez Utiliser le flux de données de modification sur Databricks.

Si vous utilisez les Lakeflow pipelines, consultez Les AUTO CDC APIs : Simplifiez la capture des données modifiées avec les pipelines.

important

Dans Databricks Runtime 12,2 LTS et versions antérieures, vous ne pouvez pas stream à partir du flux de données de modification pour une table Delta Lake avec le mappage de colonnes activé qui a subi une évolution de schéma non additive, telle que le renommage ou la suppression de colonnes. Consultez l'article sur le mappage de colonnes et le streaming.

Utiliser des vues matérialisées

Les vues matérialisées gèrent automatiquement les modifications en amont en recalculant les résultats lorsque les données sources changent. Si vous n'avez pas besoin de la latence la plus faible possible et que vous souhaitez éviter de gérer la complexité du streaming, une vue matérialisée peut simplifier votre architecture. Les vues matérialisées sont disponibles dans les LakeFlow Pipelines et les pipelines autonomes. Voir les vues matérialisées.

Exemple

Par exemple, supposons que vous ayez une table user_events avec les colonnes date, user_email et action qui est partitionnée par date. Vous Stream out of the user_events table et vous devez en supprimer les données en raison du GDPR.

skipChangeCommits vous permet de supprimer des données dans plusieurs partitions (dans cet exemple, en filtrant sur user_email). Utilisez la syntaxe suivante :

Scala
spark.readStream
.option("skipChangeCommits", "true")
.table("user_events")

Si vous mettez à jour un user_email avec la déclaration UPDATE, le fichier contenant le user_email en question est réécrit. Utilisez skipChangeCommits pour ignorer les fichiers de données modifiés.

Databricks recommande d'utiliser skipChangeCommits au lieu de ignoreDeletes, sauf si vous êtes certain que les suppressions sont toujours des suppressions de partition complètes.

Utilisez foreachBatch pour les écritures de table idempotentes

remarque

Databricks recommande de configurer une écriture en streaming distincte pour chaque puits que vous souhaitez mettre à jour au lieu d'utiliser foreachBatch. Les écritures vers plusieurs puits dans foreachBatch réduisent la parallélisation et augmentent la latence globale car les écritures vers plusieurs tables sont sérialisées dans foreachBatch.

Les tables Delta Lake prennent en charge les options DataFrameWriter suivantes pour rendre les écritures vers plusieurs tables au sein de foreachBatch idempotentes :

  • txnAppId: une chaîne unique que vous pouvez transmettre à chaque écriture de DataFrame. Par exemple, vous pouvez utiliser l'ID StreamingQuery comme txnAppId. txnAppId peut être toute chaîne unique générée par l'utilisateur et n'a pas besoin d'être liée à l'ID de Stream.
  • txnVersion: Un nombre croissant de manière monotone qui agit comme version de transaction.

Delta Lake utilise txnAppId et txnVersion pour identifier et ignorer les écritures en double. Par exemple, après qu'une défaillance ait interrompu une écriture par batch, vous pouvez réexécuter le batch avec les mêmes txnAppId et txnVersion pour identifier et ignorer correctement les doublons. Consultez Utiliser foreachBatch pour écrire dans des destinations de données arbitraires.

attention

Si vous supprimez le checkpoint de streaming et redémarrez la query avec un nouveau checkpoint, vous devez fournir un txnAppId différent. Les nouveaux points de contrôle start avec un ID de batch de 0. Delta Lake utilise l'ID du batch et txnAppId comme clé unique, et ignore les batchs dont les valeurs ont déjà été vues.

L'exemple de code suivant illustre ce modèle :

Python
app_id = ... # A unique string that is used as an application ID.

def writeToDeltaLakeTableIdempotent(batch_df, batch_id):
batch_df.write.format(...).option("txnVersion", batch_id).option("txnAppId", app_id).save(...) # location 1
batch_df.write.format(...).option("txnVersion", batch_id).option("txnAppId", app_id).save(...) # location 2

streamingDF.writeStream.foreachBatch(writeToDeltaLakeTableIdempotent).start()

Upsert à partir de queries de streaming en utilisant foreachBatch

Vous pouvez utiliser merge et foreachBatch pour écrire des upserts complexes à partir d'une query de streaming dans une table Delta Lake. Voir Utiliser foreachBatch pour écrire dans des destinations de données arbitraires.

Cette approche présente de nombreuses applications :

remarque
  • Vérifiez que votre instruction merge à l'intérieur de foreachBatch est idempotente. Autrement, les redémarrages de la query streaming peuvent appliquer l'opération sur le même batch de données plusieurs fois. Consultez Utiliser foreachBatch pour les écritures de table idempotentes.

  • Lorsque merge est utilisé dans foreachBatch, la métrique du débit des données d'entrée peut renvoyer un multiple du débit réel auquel les données sont générées à la source. merge lit les données d'entrée plusieurs fois, ce qui multiplie les métriques. Pour éviter la multiplication des métriques, mettez en cache le DataFrame batch avant merge, puis supprimez-le du cache après merge.

    Le taux de données d'entrée est disponible via StreamingQueryProgress et dans le graphe de taux de streaming du Notebook. Consulter monitoring des query Structured Streaming sur Databricks.

Par exemple, vous pouvez utiliser des instructions SQL MERGE dans foreachBatch:

Scala
// Function to upsert microBatchOutputDF into Delta Lake table using merge
def upsertToDelta(microBatchOutputDF: DataFrame, batchId: Long) {
// Set the dataframe to view name
microBatchOutputDF.createOrReplaceTempView("updates")

// Use the view name to apply MERGE
// NOTE: You have to use the SparkSession that has been used to define the `updates` dataframe
microBatchOutputDF.sparkSession.sql(s"""
MERGE INTO aggregates t
USING updates s
ON s.key = t.key
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
""")
}

// Write the output of a streaming aggregation query into Delta Lake table
streamingAggregatesDF.writeStream
.foreachBatch(upsertToDelta _)
.outputMode("update")
.start()

Vous pouvez également utiliser les APIs Delta Lake pour les upserts en streaming :

Scala
import io.delta.tables.*

val deltaTable = DeltaTable.forName(spark, "table_name")

// Function to upsert microBatchOutputDF into Delta Lake table using merge
def upsertToDelta(microBatchOutputDF: DataFrame, batchId: Long) {
deltaTable.as("t")
.merge(
microBatchOutputDF.as("s"),
"s.key = t.key")
.whenMatched().updateAll()
.whenNotMatched().insertAll()
.execute()
}

// Write the output of a streaming aggregation query into Delta Lake table
streamingAggregatesDF.writeStream
.foreachBatch(upsertToDelta _)
.outputMode("update")
.start()

Définir la version initiale de la table pour traiter les changements

Par default, les streams commencent avec la dernière version de table Delta Lake disponible. Ceci inclut un instantané complet de la table à ce moment et toutes les modifications futures. Databricks vous recommande d'utiliser la version de table initiale default pour la plupart des charges de travail.

En option, vous pouvez utiliser les options suivantes pour spécifier le point de départ de la source de streaming Delta Lake sans traiter la table entière.

  • startingVersion: La version de la table Delta Lake à partir de laquelle start la lecture. Toutes les modifications de table validées à partir de la version spécifiée sont lues par le stream. Si la version spécifiée n'est pas disponible, le stream ne start pas.

    Pour trouver les versions de commit disponibles, exécutez DESCRIBE HISTORY et vérifiez le version. Pour n’afficher que les dernières modifications, spécifiez latest. Pour des informations sur les versions des tables Delta Lake, consultez Travailler avec l'historique des tables.

  • startingTimestamp: Le Timestamp de start de lecture. Toutes les modifications de table validées à partir ou après le Timestamp spécifié sont lues par le Stream. Si le Timestamp fourni précède tous les commits de table, la lecture en streaming commence avec le Timestamp le plus ancien disponible. Définir :

    • Chaîne Timestamp. Par exemple, "2019-01-01T00:00:00.000Z".
    • Une chaîne de date. Par exemple, "2019-01-01".

Vous ne pouvez pas définir startingVersion et startingTimestamp en même temps. Ces paramètres s'appliquent uniquement aux nouvelles queries de streaming. Si une query de streaming a start et que la progression a été enregistrée dans son checkpoint, ces paramètres sont ignorés.

important

Bien que vous puissiez start la source de streaming à partir d'une version ou d'un Timestamp spécifié, le schéma de la source de streaming est toujours le dernier schéma de la table Delta Lake. Vous devez vous assurer qu'il n'y a pas de modification de schéma incompatible de la table Delta Lake après la version ou le Timestamp spécifié. Sinon, la source de streaming peut renvoyer des résultats incorrects lors de la lecture des données avec un schéma incorrect.

Exemple

Par exemple, supposons que vous ayez une table user_events. Si vous souhaitez lire les modifications depuis la version 5, utilisez :

Scala
spark.readStream
.option("startingVersion", "5")
.table("user_events")

Si vous souhaitez lire les modifications depuis le 18-10-2018, utilisez :

Scala
spark.readStream
.option("startingTimestamp", "2018-10-18")
.table("user_events")

Traiter l'instantané initial sans perte de données

Cette fonctionnalité est disponible sur Databricks Runtime 11.3 LTS et versions supérieures.

Dans une query de streaming avec état et un watermark défini, le traitement des fichiers par heure de modification peut traiter les enregistrements dans le mauvais ordre. Cela peut entraîner le marquage incorrect des enregistrements par le filigrane comme des événements tardifs et leur suppression. Cela ne peut se produire que lorsque le snapshot Delta initial est traité dans l'ordre default.

Pour les Stream avec une table source Delta, la query traite d'abord toutes les données présentes dans la table et crée une version appelée l'*instantané initial*. default, les fichiers de données de la table Delta Lake sont traités en fonction du dernier fichier modifié. Cependant, l'heure de la dernière modification ne représente pas nécessairement l'ordre temporel des événements enregistrés.

Pour éviter les pertes de données lors du traitement du snapshot initial, activez l’option withEventTimeOrder. withEventTimeOrder divise la plage horaire d’événements des données de snapshot initiales en segments temporels. Chaque micro-batch traite un bucket en filtrant les données dans la plage de temps. Les options maxFilesPerTrigger et maxBytesPerTrigger sont toujours applicables pour contrôler la taille du micro-batch, mais seulement de manière approximative en raison de l’approche de traitement.

Le diagramme suivant illustre ce processus :

Instantané initial

Contraintes

  • Vous ne pouvez pas modifier withEventTimeOrder si la query du Stream a start et que l'instantané initial est en cours de traitement. Pour redémarrer avec withEventTimeOrder modifié, vous devez supprimer le checkpoint.
  • Si withEventTimeOrder est activé, vous ne pouvez pas rétrograder un Stream vers une version de Databricks Runtime qui ne prend pas en charge cette fonctionnalité tant que le traitement initial de l'instantané n'est pas terminé. Pour rétrograder, attendez la fin de l'instantané initial, ou supprimez le point de contrôle et redémarrez la query.
  • Cette fonctionnalité n'est pas prise en charge dans les scénarios suivants :
    • La colonne de temps d'événement est une colonne générée et il existe des Transformations de non-projection entre la source Delta et le filigrane.
    • Il y a un filigrane qui a plus d'une source Delta dans la requête de Stream.

Performance

Si withEventTimeOrder est activé, les performances de traitement de l'instantané initial pourraient être plus lentes. Chaque micro-batch analyse l'instantané initial pour filtrer les données dans la plage horaire d'événements correspondante. Pour améliorer les performances de filtrage :

  • Utilisez une colonne source Delta comme heure d’événement afin que le saut de données puisse être appliqué. Voir Saut de données.
  • Partitionnez la table selon la colonne du temps de l'événement.

Utilisez Spark UI pour voir combien de fichiers Delta sont analysés pour un micro-batch spécifique.

Exemple

Supposons que vous ayez une table user_events avec une colonne event_time. Votre query de streaming est une query d'agrégation. Si vous voulez vous assurer qu'aucune donnée n'est perdue pendant le traitement initial de l'instantané, vous pouvez utiliser :

Scala
spark.readStream
.option("withEventTimeOrder", "true")
.table("user_events")
.withWatermark("event_time", "10 seconds")

Vous pouvez définir withEventTimeOrder avec une configuration Spark sur le cluster pour l’appliquer à toutes les queries en streaming : spark.databricks.delta.withEventTimeOrder.enabled true.

Limiter le taux d'entrée pour améliorer les performances de traitement

By default, Structured Streaming traite autant de fichiers que possible dans chaque micro-batch. Pour limiter la quantité de données traitées par batch et gérer l'utilisation de la mémoire, stabiliser la latence ou réduire les coûts de stockage cloud, utilisez les options suivantes :

  • maxFilesPerTrigger: Le nombre de nouveaux fichiers à prendre en compte dans chaque micro-batch. La valeur par default est 1000.
  • maxBytesPerTrigger: La quantité de données traitées dans chaque micro-batch. Cette option définit un « soft max », ce qui signifie qu’un batch traite approximativement cette quantité de données et pourrait traiter plus que la limite afin de faire avancer la query en streaming dans les cas où la plus petite unité d’entrée est supérieure à cette limite. Ceci n’est pas défini par default.

Si vous utilisez à la fois maxBytesPerTrigger et maxFilesPerTrigger, le micro-batch traite les données jusqu'à ce que la limite maxFilesPerTrigger ou maxBytesPerTrigger soit atteinte.

remarque

Par default, si logRetentionDuration nettoie les transactions dans la table source et que la query en streaming tente de traiter ces versions, la query échoue pour éviter la perte de données. Vous pouvez définir l'option failOnDataLoss sur false pour ignorer les données perdues et continuer le traitement. Consultez Configurer la conservation des données pour les queries de time travel.

Contrôlez le coût du stockage cloud

Les queries en streaming disposent de plusieurs modes de trigger disponibles qui vous permettent d'équilibrer les coûts et la latence, notamment processingTime, availableNow et realTime. Voir Contrôler le coût du stockage cloud.