Utilisez ForEachBatch pour écrire vers des récepteurs de données arbitraires dans des pipelines.

Le collecteur ForEachBatch traite un flux en une série de micro-lots. Chaque lot peut être traité dans Python avec une logique personnalisée similaire à la foreachBatch d'Apache Spark Structured Streaming. Avec le sink ForEachBatch de Lakeflow Pipelines, vous pouvez transformer, fusionner ou écrire des données en continu vers une ou plusieurs cibles qui ne prennent pas en charge nativement l’écriture en continu.

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

  • Logique personnalisée pour chaque micro-lot : ForEachBatch est un récepteur flexible de diffusion en continu. Vous pouvez appliquer des actions arbitraires (telles que la fusion dans une table externe, l’écriture dans plusieurs destinations ou l’exécution d’upserts) avec du code Python.
  • Prise en charge complète de l’actualisation : les pipelines gèrent les points de contrôle par flux. Par conséquent, les points de contrôle sont réinitialisés automatiquement lorsque vous effectuez une actualisation complète de votre pipeline. Avec le récepteur ForEachBatch, vous êtes responsable de la gestion de la réinitialisation des données en aval lorsque cela se produit.
  • Prise en charge de Unity Catalog : le récepteur ForEachBatch prend en charge toutes les fonctionnalités de Unity Catalog, telles que la lecture ou l’écriture dans des volumes ou des tables de Unity Catalog.
  • Tâches de nettoyage limitées : le pipeline n'effectue pas le suivi des données écrites à partir d’un récepteur ForEachBatch. Il ne peut donc pas nettoyer ces données. Vous êtes responsable de toute gestion des données en aval.
  • Entrées du journal des événements : le journal des événements du pipeline enregistre la création et l’utilisation de chaque puits ForEachBatch. Si votre fonction Python n’est pas sérialisable, vous voyez une entrée d’avertissement dans le journal des événements avec des suggestions supplémentaires.

Note

  • Le récepteur ForEachBatch est conçu pour les requêtes de diffusion en continu, telles que append_flow. Il n'est pas destiné aux pipelines uniquement par lots ni à la sémantique AutoCDC.
  • Le récepteur ForEachBatch qui est décrit sur cette page concerne les pipelines. Apache Spark Structured Streaming prend également en charge foreachBatch. Pour plus d’informations sur Structured Streaming foreachBatch, consultez Utiliser foreachBatch pour écrire dans des récepteurs de données arbitraires.

Quand utiliser un collecteur ForEachBatch

Utilisez un récepteur ForEachBatch chaque fois que votre pipeline nécessite des fonctionnalités qui ne sont pas disponibles via un format récepteur intégré tel que delta, ou kafka. Les cas d’usage classiques sont les suivants :

  • Fusion ou opération d’upsert dans une table Delta Lake : exécutez une logique de fusion personnalisée pour chaque micro-lot (par exemple, la gestion des enregistrements mis à jour).
  • Écriture vers plusieurs destinations ou vers des destinations non prises en charge : écrivez la sortie de chaque lot dans plusieurs tables ou systèmes de stockage externes qui ne prennent pas en charge les écritures de diffusion en continu (comme certains récepteurs JDBC).
  • Applying custom logic or transformations : Manipuler des données dans Python directement (par exemple, à l’aide de bibliothèques spécialisées ou de transformations avancées).

Pour plus d’informations sur les récepteurs intégrés ou sur la création de récepteurs personnalisés avec Python, consultez Récepteurs dans les pipelines Lakeflow.

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

Actualisation complète

Étant donné que ForEachBatch utilise une requête de diffusion en continu, le pipeline effectue le suivi du répertoire de point de contrôle pour chaque flux. Lors de l’actualisation complète :

  • Le répertoire de point de contrôle est réinitialisé.
  • Votre fonction récepteur (foreach_batch_sink UDF) voit un tout nouveau cycle batch_id démarrant à 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 partant de zéro, vous devez supprimer ou tronquer manuellement les tables externes ou les emplacements que votre récepteur ForEachBatch remplit.

Utilisation des fonctionnalités du catalogue Unity

Toutes les fonctionnalités existantes du catalogue Unity dans Spark Structured Streaming foreach_batch_sink restent disponibles.

Cela inclut l’écriture dans des tables du catalogue Unity, qu’elles soient gérées ou externes. Vous pouvez écrire des micro-lots dans des tables managées ou externes de Unity Catalog exactement comme vous le feriez dans n'importe quel travail Apache Spark Structured Streaming.

Entrées du journal des événements

Lorsque vous créez un récepteur ForEachBatch, un événement SinkDefinition, accompagné de "format": "foreachBatch" est ajouté au journal des événements du pipeline.

Cela vous permet de suivre l’utilisation des récepteurs ForEachBatch et de voir les avertissements concernant votre récepteur.

Utilisation avec Databricks Connect

Si la fonction que vous fournissez n’est pas sérialisable (une exigence importante pour Databricks Connect), le journal des événements inclut une WARN entrée qui vous recommande de simplifier ou de refactoriser votre code si la prise en charge de Databricks Connect est requise.

Par exemple, si vous utilisez dbutils pour obtenir des paramètres dans une fonction UDF ForEachBatch, vous pouvez obtenir l’argument avant de l’utiliser dans la fonction UDF :

# 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)

Meilleures pratiques

  1. Conservez votre fonction ForEachBatch concise : évitez le threading, les dépendances de bibliothèque lourdes ou les manipulations de données en mémoire volumineuses. Une logique complexe ou avec état peut entraîner des erreurs de sérialisation ou des goulots d’étranglement des performances.
  2. Surveillez votre dossier de checkpoints : pour les requêtes en streaming, le pipeline gère les checkpoints par flux, et non par sortie. Si vous avez plusieurs flux dans votre pipeline, chaque flux a son propre répertoire de point de contrôle.
  3. Valider les dépendances externes : si vous vous appuyez sur des systèmes ou bibliothèques externes, vérifiez qu’elles sont installées sur tous les nœuds de cluster ou dans votre conteneur.
  4. N'oubliez pas Databricks Connect : si votre environnement est susceptible de passer à Databricks Connect à l'avenir, vérifiez que votre code est sérialisable et ne dépend pas de dbutils dans la foreach_batch_sink fonction UDF.

Limites

  • Aucune maintenance pour ForEachBatch : étant donné que votre code python peut écrire des données n'importe où, le pipeline ne peut ainsi pas nettoyer ou suivre ces données. Vous devez gérer vos propres stratégies de gestion ou de rétention des données pour les destinations dans lesquelles vous écrivez.
  • Métriques en micro-lots : les pipelines collectent des 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 la traçabilité du flux des données et des lignes difficile pour le système.
  • Prise en charge de l’écriture dans plusieurs destinations sans plusieurs lectures : certains clients peuvent utiliser ForEachBatch pour lire à partir d’une source une seule fois, puis écrire dans plusieurs destinations. Pour ce faire, vous devez inclure df.persist ou df.cache à l’intérieur de votre fonction ForEachBatch. À l’aide de ces options, Azure Databricks tente de lire les données une seule fois. Sans ces options, votre requête génère plusieurs lectures. Cela n’est pas inclus dans les exemples de code suivants.
  • Utilisation avec Databricks Connect : si votre pipeline s’exécute sur Databricks Connect, foreachBatch les fonctions définies par l’utilisateur (UDF) doivent être sérialisables et ne peuvent pas être utilisées dbutils. Le pipeline déclenche des avertissements s’il détecte une fonction définie par l'utilisateur (UDF) non sérialisée, mais il ne fait pas échouer le pipeline.
  • Logique non sérialisable : le code référençant des objets locaux, des classes ou des ressources non sélectionnables peut s’interrompre dans les contextes Databricks Connect. Utilisez des modules Python purs et vérifiez que les références (par exemple, dbutils) ne sont pas utilisées si Databricks Connect est obligatoire.

Examples

Exemple de syntaxe de base

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 d’exemples de données pour un pipeline simple

Cet exemple utilise l’exemple NYC Taxi. Il part du principe que l’administrateur de votre espace de travail a activé le catalogue Databricks Public Datasets. Pour le récepteur, modifiez my_catalog.my_schema pour en faire un catalogue et un schéma auxquels vous avez accès.

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 idempotentes les écritures vers les tables Delta Lake. Pour plus d’informations, consultez Utiliser foreachBatch pour les écritures de tables idempotentes.

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

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 de spark.sql()

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

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

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 envisagez d’exécuter votre pipeline dans un environnement non Databricks Connect, dbutils cela peut fonctionner. Toutefois, si vous utilisez Databricks Connect, dbutils n’est pas accessible dans votre foreachBatch fonction. Le pipeline peut déclencher des avertissements s’il détecte l’utilisation de dbutils pour vous aider à éviter les ruptures.

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

Yes. Vous pouvez définir plusieurs flux (avec @dp.append_flow) qui ont tous pour cible le même nom de destination, mais chacun 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 puits de données ForEachBatch peut écrire dans n’importe quel emplacement ou système spécifié, 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 résoudre les erreurs de sérialisation ou les échecs dans ma fonction ForEachBatch ?

Examinez les journaux d’activité de votre pilote de cluster ou les journaux d’événements de pipeline. Pour les problèmes de sérialisation liés à Spark Connect, vérifiez que votre fonction dépend uniquement des objets sérialisables Python et ne fait pas référence à des objets non autorisés (tels que les handles de fichiers ouverts ou dbutils).