Utilisez foreachBatch pour écrire dans des puits de données arbitraires
Cette page montre comment utiliser foreachBatch avec Structured Streaming pour écrire le résultat d'une query de streaming dans des sources de données qui n'ont pas de récepteur de streaming existant.
Le modèle de code streamingDF.writeStream.foreachBatch(...) vous permet d’appliquer des fonctions batch aux données de sortie de chaque micro-batch de la query de streaming. Les fonctions utilisées avec foreachBatch prennent deux paramètres :
- Un DataFrame qui contient les données de sortie d'un micro-batch.
- L'ID unique du micro-batch.
Vous devez utiliser foreachBatch pour les opérations de merge Delta Lake dans Structured Streaming. Consultez Insérer ou mettre à jour à partir de requêtes de streaming à l'aide de foreachBatch.
Appliquer des opérations DataFrame supplémentaires
De nombreuses opérations de DataFrame et de dataset ne sont pas prises en charge dans les DataFrames de streaming, car Spark ne prend pas en charge la génération de plans incrémentiels dans ces cas. Avec foreachBatch(), vous pouvez appliquer certaines de ces Opérations sur chaque sortie micro-batch. Par exemple, vous pouvez utiliser foreachBatch() et l'opération SQL MERGE INTO pour écrire la sortie des agrégations de streaming dans une table Delta Lake en mode de mise à jour. Voir plus de détails dans MERGE INTO.
foreachBatch()offre uniquement des garanties d'écriture au moins une fois. Toutefois, vous pouvez utiliser lebatchIdfourni à la fonction afin de dédupliquer la sortie et d'obtenir une garantie d'exécution exacte. Dans tous les cas, vous devrez réfléchir à la sémantique de bout en bout vous-même.foreachBatch()ne fonctionne pas avec le mode de traitement continu car il repose fondamentalement sur l'exécution de micro-batch d'une query en streaming. Si vous écrivez des données en mode continu, utilisezforeach()à la place.- Lorsque vous utilisez
foreachBatchavec un opérateur avec état, il est important de consommer entièrement chaque batch avant la fin du traitement. Consultez Consommer complètement chaque DataFrame batch
Gérer les DataFrames vides
foreachBatch() peut recevoir un DataFrame vide, et votre code doit gérer ce scénario. Autrement, votre query pourrait échouer.
Par exemple, lorsque Delta Lake est la source de streaming, ces scénarios peuvent passer un DataFrame vide à foreachBatch():
OPTIMIZEsans fichiers à traiter : Lorsqu’uneOPTIMIZEopération s’exécute sur la table source Delta Lake mais qu’il n’y a aucun fichier à traiter, Structured Streaming écrit une entrée de Log de décalage pour incrémenter la version de la table. Cela produit un micro-batch vide sur le récepteur même si aucun fichier n’est lu.- Élagage de fichiers au niveau du plan physique : si le pushdown de prédicat ou l'élagage de fichiers élimine tous les enregistrements au niveau du plan physique, le résultat est un commit vide vers le récepteur.
Le code utilisateur doit gérer les DataFrames vides pour permettre une bonne opération. Veuillez consulter les exemples ci-dessous :
- Python
- Scala
def process_batch(output_df, batch_id):
# Process valid DataFrames only
if not output_df.isEmpty():
# business logic
pass
streamingDF.writeStream.foreachBatch(process_batch).start()
.foreachBatch(
(outputDf: DataFrame, bid: Long) => {
// Process valid DataFrames only
if (!outputDf.isEmpty) {
// business logic
}
}
).start()
Modifications de comportement pour foreachBatch dans Databricks Runtime 14.0
Dans Databricks Runtime 14.0 et versions ultérieures sur un compute configuré avec le mode d'accès standard, les changements de comportement suivants s'appliquent :
print()Les commandes écrivent la sortie dans les Logs Driver.- Vous ne pouvez pas accéder au sous-module
dbutils.widgetsà l'intérieur de la fonction. - Tous les fichiers, modules ou objets référencés dans la fonction doivent être sérialisables et disponibles sur Spark.
Réutiliser les sources de données batch existantes
En utilisant foreachBatch(), vous pouvez utiliser les enregistreurs de données batch existants pour les puits de données qui pourraient ne pas prendre en charge le Structured Streaming. Voici quelques exemples :
De nombreuses autres sources de données batch peuvent être utilisées à partir de foreachBatch(). Consultez Connecter aux sources de données et aux services externes.
Écriture vers plusieurs emplacements
Si vous devez écrire le résultat d'une requête de streaming vers plusieurs emplacements, Databricks recommande d'utiliser plusieurs rédacteurs Structured Streaming pour une meilleure parallélisation et un meilleur throughput.
L'utilisation de foreachBatch pour écrire vers plusieurs récepteurs sérialise l'exécution des écritures de streaming, ce qui peut augmenter la latence pour chaque micro-batch.
Si vous utilisez foreachBatch pour écrire dans plusieurs tables Delta Lake, consultez Utiliser foreachBatch pour les écritures de tables idempotentes.
Consommez complètement chaque DataFrame de batch
Lorsque vous utilisez des opérateurs avec état (par exemple, en utilisant dropDuplicatesWithinWatermark), chaque itération de batch doit consommer l'intégralité du DataFrame ou redémarrer la query. Si vous ne consommez pas l'intégralité du DataFrame, la query streaming échouera avec le prochain batch.
Cela peut se produire dans plusieurs cas. Les exemples suivants montrent comment corriger les query qui ne consomment pas correctement un DataFrame.
Utilisation intentionnelle d'un sous-ensemble du batch
Si vous ne vous intéressez qu'à un sous-ensemble du batch, vous pourriez avoir un code tel que le suivant.
from pyspark.sql.functions import expr
stream = spark.readStream.format("rate").option("rowsPerSecond", "100").load()
# creates a stateful operator:
streamWithWatermark = stream.withWatermark("timestamp", "15 minutes").dropDuplicatesWithinWatermark()
def partial_func(batch_df, batch_id):
batch_df.show(2)
q = streamWithWatermark.writeStream \
.foreachBatch(partial_func) \
.option("checkpointLocation", checkpoint_dir) \
.trigger(processingTime='2 seconds') \
.start()
Dans ce cas, le batch_df.show(2) ne gère que les deux premiers éléments du batch, ce qui est attendu, mais s'il y a plus d'éléments, ils doivent être consommés. Le code suivant consomme le DataFrame complet.
from pyspark.sql.functions import expr
stream = spark.readStream.format("rate").option("rowsPerSecond", "100").load()
streamWithWatermark = stream.withWatermark("timestamp", "15 minutes").dropDuplicatesWithinWatermark()
# function to do nothing with a row
def do_nothing(row):
pass
def partial_func(batch_df, batch_id):
batch_df.show(2)
batch_df.foreach(do_nothing) # silently consume the rest of the batch
q = streamWithWatermark.writeStream \
.foreachBatch(partial_func) \
.option("checkpointLocation", checkpoint_dir) \
.trigger(processingTime='2 seconds') \
.start()
Ici, la fonction do_nothing ignore silencieusement le reste du DataFrame.
Gestion d'une erreur dans un batch
Pour la gestion des erreurs dans foreachBatch, Databricks vous recommande de laisser la query de streaming échouer rapidement et de vous fier plutôt à la couche d'orchestration, telle que Lakeflow Jobs ou Apache Airflow, pour gérer la logique de nouvelle tentative. C'est beaucoup plus sûr que de créer des boucles de nouvelle tentative complexes dans votre code, où des pertes de données peuvent se produire.
Voici les directives basées sur votre objectif de rédaction :
Cible | Exemples | Guides |
|---|---|---|
Opérations DataFrame | Tables Delta Lake | Vous devez utiliser les options d'écriture |
Code personnalisé et destinations externes |
| Implémentez votre propre idempotence. Vous devez partir du principe que toute opération peut et sera sujette à de nouvelles tentatives sur l'ensemble des batches. Si le |
Voici quelques exemples de types d'exceptions et de recommandations sur la manière de les gérer dans foreachBatch:
Type d'exception | Exemples | Action recommandée |
|---|---|---|
Erreurs de puits transitoires |
| Catch : réessayer ou envoyer à une file d'attente de lettres mortes |
Violations de contrainte de clé ou de doublons lorsque le récepteur est idempotent. |
| Catch : Logs et suppression |
Erreurs personnalisées relançables | Exceptions de socket encapsulées, erreurs de base de données récupérables | Catch : incrémentez les métriques et permettez une poursuite contrôlée |
Erreurs de logique ou de schéma |
| Propagate : laissez Spark échouer la query |
Erreurs de récepteur non récupérables ou bogues logiques non interceptés |
| Propagate : laissez Spark échouer la query |
Échecs critiques |
| Propagate : laissez Spark échouer la query |
Exemples de code : gestion des exceptions
Les exemples suivants provoquent intentionnellement une erreur dans foreach pour présenter différentes approches de gestion de l'erreur :
from pyspark.sql.functions import expr
stream = spark.readStream.format("rate").option("rowsPerSecond", "100").load()
# creates a stateful operator:
streamWithWatermark = stream.withWatermark("timestamp", "15 minutes").dropDuplicatesWithinWatermark()
def foreach_func(row):
# handle the row, but in this case, for the sample, will just raise an error:
raise Exception('error')
def partial_func(batch_df, batch_id):
try:
batch_df.foreach(foreach_func)
except Exception as e:
print(e) # or whatever error handling you want to have
q = streamWithWatermark.writeStream \
.foreachBatch(partial_func) \
.option("checkpointLocation", checkpoint_dir) \
.trigger(processingTime='2 seconds') \
.start()
Le code ci-dessus gère et supprime silencieusement l'erreur, et peut ne pas consommer le reste du batch. Il existe deux options pour gérer cette situation.
Tout d'abord, vous pouvez relancer l'erreur, qui la transmet à votre couche d'orchestration pour relancer le batch. Cela peut résoudre l'erreur, s'il s'agit d'un problème transitoire, ou la transmettre à votre équipe d'Opérations pour tenter de la corriger manuellement. Pour ce faire, modifiez le code partial_func pour qu'il ressemble à ceci :
def partial_func(batch_df, batch_id):
try:
batch_df.foreach(foreach_func)
except Exception as e:
print(e) # or whatever error handling you want to have
raise e # re-raise the issue
Deuxièmement, si vous souhaitez intercepter l'exception et ignorer le reste du batch, vous pouvez modifier le code pour utiliser la fonction do_nothing afin d'ignorer discrètement le reste du batch.
from pyspark.sql.functions import expr
stream = spark.readStream.format("rate").option("rowsPerSecond", "100").load()
# creates a stateful operator:
streamWithWatermark = stream.withWatermark("timestamp", "15 minutes").dropDuplicatesWithinWatermark()
def foreach_func(row):
# handle the row, but in this case, for the sample, will just raise an error:
raise Exception('error')
# function to do nothing with a row
def do_nothing(row):
pass
def partial_func(batch_df, batch_id):
try:
batch_df.foreach(foreach_func)
except Exception as e:
print(e) # or whatever error handling you want to have
batch_df.foreach(do_nothing) # silently consume the remainder of the batch
q = streamWithWatermark.writeStream \
.foreachBatch(partial_func) \
.option("checkpointLocation", checkpoint_dir) \
.trigger(processingTime='2 seconds') \
.start()
Écrire les enregistrements en échec dans une file d'attente de lettres mortes
Par default, les queries échouent immédiatement lorsque de mauvais enregistrements arrivent. Évitez ces interruptions en configurant une table Delta Lake secondaire en tant que file d'attente de messages non distribuables (DLQ).
Avec une DLQ, le système achemine les enregistrements en échec vers la table secondaire, et continue de traiter les données valides sans interruption. La table DLQ vous permet d'inspecter et de retraiter les enregistrements erronés ultérieurement.
Utilisez cette méthode lorsque :
- Votre Stream contient des données variées qui pourraient violer les contraintes de schéma ou les règles métier.
- Les règles d'audit ou de conformité exigent que vous conserviez tous les enregistrements.
Exemple
L’exemple suivant utilise foreachBatch pour diviser chaque micro-batch en enregistrements valides et non valides, puis écrit chaque sous-ensemble dans sa propre table Delta Lake.
- Python
- Scala
from pyspark.sql.functions import current_timestamp, lit
main_table = "catalog.schema.orders"
dlq_table = "catalog.schema.orders_dlq"
app_id = "orders-streaming-job"
def process_orders(batch_df, batch_id):
if batch_df.isEmpty():
return
valid_condition = "order_amount > 0 AND customer_id IS NOT NULL"
# Write valid records to the main table with idempotency options
batch_df.filter(valid_condition).write \
.format("delta") \
.mode("append") \
.option("txnVersion", batch_id) \
.option("txnAppId", app_id) \
.saveAsTable(main_table)
# Route invalid records to the dead-letter queue
invalid_df = batch_df.filter(f"NOT ({valid_condition})")
if not invalid_df.isEmpty():
invalid_df \
.withColumn("dlq_batch_id", lit(batch_id)) \
.withColumn("dlq_ingest_time", current_timestamp()) \
.write \
.format("delta") \
.mode("append") \
.option("txnVersion", batch_id) \
.option("txnAppId", app_id) \
.saveAsTable(dlq_table)
spark.readStream \
.format("delta") \
.table("catalog.schema.raw_orders") \
.writeStream \
.foreachBatch(process_orders) \
.option("checkpointLocation", "/path/to/checkpoint") \
.start()
import org.apache.spark.sql.DataFrame
import org.apache.spark.sql.functions.{current_timestamp, lit}
val mainTable = "catalog.schema.orders"
val dlqTable = "catalog.schema.orders_dlq"
val appId = "orders-streaming-job"
def processOrders(batchDf: DataFrame, batchId: Long): Unit = {
if (batchDf.isEmpty) return
val validCondition = "order_amount > 0 AND customer_id IS NOT NULL"
// Write valid records to the main table with idempotency options
batchDf.filter(validCondition).write
.format("delta")
.mode("append")
.option("txnVersion", batchId)
.option("txnAppId", appId)
.saveAsTable(mainTable)
// Route invalid records to the dead-letter queue
val invalidDf = batchDf.filter(s"NOT ($validCondition)")
if (!invalidDf.isEmpty) {
invalidDf
.withColumn("dlq_batch_id", lit(batchId))
.withColumn("dlq_ingest_time", current_timestamp())
.write
.format("delta")
.mode("append")
.option("txnVersion", batchId)
.option("txnAppId", appId)
.saveAsTable(dlqTable)
}
}
spark.readStream
.format("delta")
.table("catalog.schema.raw_orders")
.writeStream
.foreachBatch(processOrders _)
.option("checkpointLocation", "/path/to/checkpoint")
.start()
Les deux écritures utilisent txnVersion et txnAppId pour garantir l'idempotence. Si Spark rejoue un batch avec le même batchId, Delta Lake ignore l'écriture en double. Consultez Utilisez foreachBatch pour les écritures de table idempotentes.