Aller au contenu principal

Utilisez ForEachBatch pour écrire dans des puits de données arbitraires dans les pipelines

Le puits ForEachBatch traite un Stream comme une série de micro-batchs. Chaque batch peut être traité en Python avec une logique personnalisée similaire à foreachBatch de Apache Spark Structured Streaming. Avec le récepteur ForEachBatch des Lakeflow Pipelines, vous pouvez transformer, Merge, ou écrire des données streaming vers une ou plusieurs cibles qui ne prennent pas en charge les écritures en streaming en mode natif.

Le récepteur ForEachBatch offre les fonctionnalités suivantes :

  • **Logique personnalisée pour chaque micro-lot** : ForEachBatch est un récepteur de streaming flexible. Vous pouvez appliquer des actions arbitraires (telles que la fusion dans une table externe, l'écriture vers plusieurs destinations ou la réalisation d'upserts) avec du code Python.
  • Prise en charge de la complète refresh : les pipelines gèrent les points de contrôle par flux, de sorte que les points de contrôle se Reset automatiquement lorsque vous effectuez une complète refresh de votre pipeline. Avec le récepteur ForEachBatch, vous êtes responsable de la gestion de la Reset des données en aval lorsque cela se produit.
  • Prise en charge d'Unity Catalog : ForEachBatch sink prend en charge toutes les fonctionnalités d'Unity Catalog, telles que la lecture et l'écriture dans des volumes ou des tables Unity Catalog.
  • Nettoyage limité : le pipeline ne suit pas les données écrites à partir d'un récepteur ForEachBatch, il ne peut donc pas nettoyer ces données. Vous êtes responsable de toute gestion de données en aval.
  • Entrées du log des événements : le pipeline log des événements du pipeline enregistre la création et l'utilisation de chaque sink ForEachBatch. Si votre fonction Python n'est pas sérialisable, vous verrez une entrée d'avertissement dans le journal des événements avec des suggestions supplémentaires.
remarque
  • Le récepteur ForEachBatch est conçu pour les query streaming, telles que append_flow. Il n'est pas destiné aux pipelines en mode batch uniquement ou à la sémantique AutoCDC.
  • Le récepteur ForEachBatch décrit sur cette page est destiné aux pipelines. Apache Spark Structured Streaming prend également en charge foreachBatch. Pour plus d'informations sur le Structured Streaming foreachBatch, consultez Utiliser foreachBatch pour écrire dans des récepteurs de données arbitraires.

Quand utiliser un récepteur ForEachBatch

Utilisez un sink ForEachBatch chaque fois que votre pipeline nécessite une fonctionnalité qui n'est pas disponible via un format de sink intégré tel que delta ou kafka. Les cas d'utilisation typiques incluent :

  • Fusion ou upsert dans une table Delta Lake : exécutez une logique de Merge personnalisée pour chaque micro-batch (par exemple, le traitement des enregistrements mis à jour).
  • Écriture vers des destinations multiples ou non prises en charge : écrivez le résultat de chaque batch dans plusieurs tables ou systèmes de stockage externes qui ne prennent pas en charge les écritures en streaming (comme certains sinks JDBC).
  • Application d’une logique personnalisée ou de Transformations : manipulez les données directement dans Python (par exemple, en utilisant des bibliothèques spécialisées ou des Transformations avancées).

Pour des informations sur les puits intégrés, ou la création de puits personnalisés avec Python, veuillez consulter Puits dans Lakeflow pipelines.

Pour la référence de l'API Python @dp.foreach_batch_sink(), consultez foreach_batch_sink.

Full refresh complète

Étant donné que ForEachBatch utilise une query en streaming, le pipeline suit le répertoire de point de contrôle pour chaque flux. Sur refresh :

  • Le répertoire du point de contrôle est reset.
  • Votre fonction de sink (foreach_batch_sink UDF) voit un tout nouveau cycle batch_id à partir de 0.
  • Les données de votre système cible ne sont pas automatiquement nettoyées par le pipeline (car le pipeline ne sait pas où vos données sont écrites). Si vous avez besoin d'un scénario vierge, vous devez manuellement supprimer ou tronquer les tables ou emplacements externes que votre récepteur ForEachBatch remplit.

Utilisation des fonctionnalités de Unity Catalog

Toutes les fonctionnalités existantes d'Unity Catalog dans Spark Structured Streaming foreach_batch_sink restent disponibles.

Cela inclut l’écriture dans des tables Unity Catalog gérées ou externes. Vous pouvez écrire des micro-batchs dans des tables Unity Catalog gérées ou externes exactement comme vous le feriez dans n’importe quel job Apache Spark Structured Streaming.

Entrées du Logs d'événements

Lorsque vous créez un sink ForEachBatch, un événement SinkDefinition, avec "format": "foreachBatch", est ajouté au log des événements du pipeline.

Ceci vous permet de suivre l'utilisation des puits ForEachBatch et de voir les avertissements concernant votre puits.

Utilisation avec Databricks Connect

Si la fonction que vous fournissez n'est **pas sérialisable** (une exigence importante pour Databricks Connect), le Log d'événements inclut une WARN entrée recommandant de simplifier ou de refactoriser votre code si le support de Databricks Connect est requis.

Par exemple, si vous utilisez dbutils pour obtenir des parameters dans un UDF ForEachBatch, vous pouvez plutôt obtenir l'argument avant de l'utiliser dans l'UDF :

Python
# Instead of accessing parameters within the UDF...
def foreach_batch(df, batchId):
value = dbutils.widgets.get ("X") + str (i)

# ...get the parameters first, and use them within the UDF:
argX = dbutils.widgets.get ("X")

def foreach_batch(df, batchId):
value = argX + str (i)

Bonnes pratiques

  1. Gardez votre fonction ForEachBatch concise : Évitez le threading, les dépendances de bibliothèque lourdes ou les manipulations de données importantes en mémoire. Une logique complexe ou avec état peut entraîner des erreurs de sérialisation ou des goulots d'activité de performance.
  2. Surveillez votre dossier de point de contrôle : pour les requêtes de streaming, le pipeline gère les points de contrôle par flux, et non par récepteur. Si votre pipeline comporte plusieurs flux, chaque flux possède son propre répertoire de point de contrôle.
  3. Validez les dépendances externes : si vous comptez sur des systèmes ou des bibliothèques externes, vérifiez qu'ils sont installés sur tous les nœuds de cluster ou dans votre conteneur.
  4. Soyez attentif à Databricks Connect : Si votre environnement est susceptible de passer à Databricks Connect à l'avenir, assurez-vous que votre code est sérialisable et ne repose pas sur dbutils au sein de l'UDF foreach_batch_sink.

Limitations

  • **Pas de nettoyage pour ForEachBatch** : parce que votre code Python personnalisé peut écrire des données n’importe où, le pipeline ne peut pas nettoyer ou suivre ces données. Vous devez gérer vos propres politiques de gestion des données ou de rétention pour les destinations vers lesquelles vous écrivez.
  • Métriques en micro-batch : Les pipelines collectent les métriques de streaming, mais certains scénarios peuvent entraîner des métriques incomplètes ou inhabituelles lors de l'utilisation de ForEachBatch. Cela est dû à la flexibilité sous-jacente de ForEachBatch qui rend le suivi des flux de données et des lignes difficile pour le système.
  • Prise en charge de l’écriture vers plusieurs destinations sans lectures multiples : Certains clients peuvent utiliser ForEachBatch pour lire à partir d’une source une seule fois, puis écrire vers plusieurs destinations. Pour ce faire, vous devez inclure df.persist ou df.cache dans votre fonction ForEachBatch. Avec ces options, Databricks tente de lire les données une seule fois. Sans ces options, votre query entraîne plusieurs lectures. Ceci n'est pas inclus dans les exemples de code suivants.
  • Utilisation avec Databricks Connect : Si votre pipeline s'exécute sur Databricks Connect, les foreachBatch fonctions définies par l'utilisateur (UDF) doivent être sérialisables et ne peuvent pas utiliser dbutils. Le pipeline génère des avertissements s'il détecte une UDF non sérialisable, mais ne fait pas échouer le pipeline.
  • Logique non sérialisable : Le code qui référence des objets locaux, des classes ou des Ressources non sérialisables peut échouer dans les contextes Databricks Connect. Utilisez des modules Python purs et confirmez que les références (par exemple, dbutils) ne sont pas utilisées si Databricks Connect est une exigence.

Exemples

Exemple de syntaxe de base

Python
from pyspark import pipelines as dp

# Create a ForEachBatch sink
@dp.foreach_batch_sink(name = "my_foreachbatch_sink")
def feb_sink(df, batch_id):
# Custom logic here. You can perform merges,
# write to multiple destinations, etc.
return

# Create source data for example:
@dp.table()
def example_source_data():
return spark.range(5)

# Add sink to an append flow:
@dp.append_flow(
target="my_foreachbatch_sink",
)
def my_flow():
return spark.readStream.format("delta").table("example_source_data")

Utilisation de données d'exemple pour un pipeline simple

Cet exemple utilise l'échantillon NYC Taxi. Il est supposé que votre administrateur de Workspace a activé le catalogue Databricks Public Datasets. Pour le récepteur, modifiez my_catalog.my_schema pour un catalogue et un schéma auxquels vous avez accès.

Python
from pyspark import pipelines as dp
from pyspark.sql.functions import current_timestamp

# Create foreachBatch sink
@dp.foreach_batch_sink(name = "my_foreach_sink")
def my_foreach_sink(df, batch_id):
# Custom logic here. You can perform merges,
# write to multiple destinations, etc.
# For this example, we are adding a timestamp column.
enriched = df.withColumn("processed_timestamp", current_timestamp())
# Write to a Delta location
enriched.write \
.format("delta") \
.mode("append") \
.saveAsTable("my_catalog.my_schema.trips_sink_delta")
# Return is optional here, but generally not used for the sink
return

# Create an append flow that reads sample data,
# and sends it to the ForEachBatch sink
@dp.append_flow(
target="my_foreach_sink",
)
def taxi_source():
df = spark.readStream.table("samples.nyctaxi.trips")
return df

Écriture vers plusieurs destinations

Cet exemple écrit vers plusieurs destinations. Il démontre l'utilisation de txnVersion et txnAppId pour rendre les écritures dans les tables Delta Lake idempotentes. Pour plus de détails, consultez Utilisez foreachBatch pour les écritures de tables idempotentes.

Supposons que nous écrivons dans deux tables, table_a et table_b, et que dans un batch, l'écriture dans table_a réussit tandis que l'écriture dans table_b échoue. Lorsque le batch est réexécuté, la paire (txnVersion, txnAppId) permettra à Delta d'ignorer l'écriture en double dans table_a, et d'écrire uniquement le batch dans table_b.

Python
from pyspark import pipelines as dp

app_id = "my-app-name" # different applications that write to the same table should have unique txnAppId

# Create the ForEachBatch sink
@dp.foreach_batch_sink(name="user_events_feb")
def user_events_handler(df, batch_id):
# Optionally do transformations, logging, or merging logic
# ...

# Write to a Delta table
df.write \
.format("delta") \
.mode("append") \
.option("txnVersion", batch_id) \
.option("txnAppId", app_id) \
.saveAsTable("my_catalog.my_schema.example_table_1")

# Also write to a JSON file location
df.write \
.format("json") \
.mode("append") \
.option("txnVersion", batch_id) \
.option("txnAppId", app_id) \
.save("/tmp/json_target")
return

# Create source data for example
@dp.table()
def example_source():
return spark.range(5)


# Create the append flow, and target the ForEachBatch sink
@dp.append_flow(target="user_events_feb", name="user_events_flow")
def read_user_events():
return spark.readStream.format("delta").table("example_source")

Utilisation spark.sql()

Vous pouvez utiliser spark.sql() dans votre récepteur ForEachBatch, comme dans l'exemple suivant.

Python
from pyspark import pipelines as dp
from pyspark.sql import Row

@dp.foreach_batch_sink(name = "example_sink")
def feb_sink(df, batch_id):
df.createOrReplaceTempView("df_view")
df.sparkSession.sql("MERGE INTO target_table AS tgt " +
"USING df_view AS src ON tgt.id = src.id " +
"WHEN MATCHED THEN UPDATE SET tgt.id = src.id * 10 " +
"WHEN NOT MATCHED THEN INSERT (id) VALUES (id)"
)
return

# Create target delta table
spark.range(5).write.format("delta").mode("overwrite").saveAsTable("target_table")

# Create source table
@dp.table()
def src_table():
return spark.range(5)

@dp.append_flow(
target="example_sink",
)
def example_flow():
return spark.readStream.format("delta").table("source_table")

Fusion avec une table Delta Lake externe

Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col
from delta.tables import DeltaTable

@dp.foreach_batch_sink(name = "external_merge_feb")
def foreachBatchFunc(df, batchId):
out = DeltaTable.forName(df.sparkSession, $table)
out.alias("target") \
.merge(df.alias("source"), "source.value = target.value") \
.whenMatchedUpdateAll() \
.whenNotMatchedInsertAll() \
.whenNotMatchedBySourceDelete() \
.execute()

@dp.update_flow(
target="external_merge_feb",
name="merge_flow"
)
def read_data():
return (
spark.readStream.format("delta")
.load("/tmp/source_delta_table")
.filter(col("value").isNotNull())
)

Questions fréquemment posées (FAQ)

Puis-je utiliser dbutils dans mon récepteur ForEachBatch ?

Si vous prévoyez d'exécuter votre pipeline dans un environnement non Databricks Connect, dbutils peut fonctionner. Cependant, si vous utilisez Databricks Connect, dbutils n'est pas accessible au sein de votre fonction foreachBatch. Le pipeline peut générer des avertissements s’il détecte une utilisation de dbutils pour vous aider à éviter les interruptions.

Puis-je utiliser plusieurs flux avec un seul récepteur ForEachBatch ?

Oui. Vous pouvez définir plusieurs flux (avec @dp.append_flow) qui ciblent tous le même nom de récepteur, mais chacun d'entre eux conserve ses propres points de contrôle.

Le pipeline gère-t-il la conservation ou le nettoyage des données pour ma cible ?

Non. Étant donné que le récepteur ForEachBatch peut écrire à n'importe quel emplacement ou système arbitraire, le pipeline ne peut pas gérer ou supprimer automatiquement les données dans cette cible. Vous devez gérer ces opérations dans le cadre de votre code personnalisé ou de vos processus externes.

Comment dépanner les erreurs de sérialisation ou les échecs dans ma fonction ForEachBatch ?

Consultez vos logs de Driver de clusters ou logs d'événements de pipeline. Pour les problèmes de sérialisation liés à Spark Connect, vérifiez que votre fonction ne dépend que d'objets Python sérialisables et ne référence pas d'objets non autorisés (tels que des descripteurs de fichiers ouverts ou dbutils).