Aller au contenu principal

Utilisez des puits dans les pipelines.

Utilisez l'API sink du pipeline Lakeflow avec des flux pour écrire les enregistrements transformés par un pipeline vers une destination de données externe. Les récepteurs de données externes incluent les tables gérées et externes du Unity Catalog, et les services de streaming d'événements tels qu'Apache Kafka ou Azure Event Hub. Vous pouvez également utiliser des destinations de données pour écrire dans des sources de données personnalisées en écrivant du code Python pour cette source de données.

Pour un aperçu des concepts de récepteur et quand les utiliser, consultez Récepteurs dans Lakeflow Pipelines.

remarque

Workflow de puits

Lorsque les données d'événement sont ingérées d'une source de streaming dans votre pipeline, vous traitez et affinez ces données dans des transformations de votre pipeline. Vous utilisez ensuite le traitement de flux d'ajout pour Stream les enregistrements de données transformés vers un puits. Vous créez ce récepteur à l'aide de la fonction create_sink(). Pour plus de détails sur la fonction create_sink, consultez la référence de l'API sink.

Si vous avez un pipeline qui crée ou traite vos données d'événements de streaming et prépare les enregistrements de données pour l'écriture, alors vous êtes prêt à utiliser un sink.

L'implémentation d'un récepteur comprend deux étapes :

  1. Créer le puits.
  2. Utilisez un flux d'ajout ou un flux de mise à jour pour écrire les enregistrements préparés dans le récepteur.

Créer un puits

Databricks prend en charge plusieurs types de destinations dans lesquelles vous écrivez vos enregistrements traités à partir de vos données de Stream :

  • Puits de tables Delta (y compris les tables gérées et externes de Unity Catalog)
  • Puits Apache Kafka
  • Destinations Azure Event Hub
  • Sinks personnalisés écrits en Python, à l'aide des sources de données personnalisées Python

Vous trouverez ci-dessous des exemples de configurations pour les récepteurs Delta, Kafka et Azure Event Hub, ainsi que des sources de données personnalisées Python :

Pour créer un récepteur Delta par chemin de fichier :

Python
dp.create_sink(
name = "delta_sink",
format = "delta",
options = {"path": "/Volumes/catalog_name/schema_name/volume_name/path/to/data"}
)

Pour créer un récepteur Delta par nom de table à l'aide d'un chemin de catalogue et de schéma entièrement qualifié :

Python
dp.create_sink(
name = "delta_sink",
format = "delta",
options = { "tableName": "catalog_name.schema_name.table_name" }
)

Pour plus de détails sur l'utilisation de la fonction create_sink, consultez la référence de l'API du sink.

Une fois votre sink créé, vous pouvez commencer à diffuser en streaming les enregistrements traités vers le sink.

Écrire vers un récepteur avec un flux d'ajout

Une fois votre sink créé, la prochaine étape consiste à y écrire les enregistrements traités en le spécifiant comme cible pour les enregistrements générés par un flux d'ajout. Pour ce faire, vous spécifiez votre sink comme valeur target dans le décorateur append_flow.

  • Pour les tables gérées et externes de Unity Catalog, utilisez le format delta et spécifiez le chemin ou le nom de la table dans les options. Votre pipeline doit être configuré pour utiliser Unity Catalog.
  • Pour les sujets Apache Kafka, utilisez le format kafka et spécifiez le nom du sujet, les informations de connexion et les informations d'authentification dans les options. Ce sont les mêmes options qu'un récepteur Kafka Spark Structured Streaming prend en charge. Voir Configurer l'enregistreur Kafka Structured Streaming.
  • Pour Azure Event Hubs, utilisez le format kafka et spécifiez le nom d'Event Hubs, les informations de connexion et les informations d'authentification dans les options. Ce sont les mêmes options prises en charge dans un récepteur Spark Structured Streaming Event Hubs qui utilise l'interface Kafka. Voir Authentification.

Vous trouverez ci-dessous des exemples de configuration de flux pour écrire vers des récepteurs Delta, Kafka et Azure Event Hub avec des enregistrements traités par votre pipeline.

Python
@dp.append_flow(name = "delta_sink_flow", target="delta_sink")
def delta_sink_flow():
return(
spark.readStream.table("spark_referrers")
.selectExpr("current_page_id", "referrer", "current_page_title", "click_count")
)

Pour plus de détails sur le décorateur append_flow, consultez default flows et flux d'ajout.

Limitations

  • Seule l'API Python est prise en charge. SQL n'est pas pris en charge.

  • Seules les queries de streaming sont prises en charge. Les requêtes batch ne sont pas prises en charge.

  • Seuls append_flow et update_flow peuvent être utilisés pour écrire dans des récepteurs. Les autres flux, tels que create_auto_cdc_flow, ne sont pas pris en charge, et vous ne pouvez pas utiliser un récepteur dans une définition de dataset de pipeline. Par exemple, ce qui suit n'est pas pris en charge :

    Python
    @table("from_sink_table")
    def fromSink():
    return read_stream("my_sink")
  • Pour les sinks Delta, le nom de la table doit être entièrement qualifié. Plus précisément, pour les tables externes gérées par Unity Catalog, le nom de la table doit être de la forme <catalog>.<schema>.<table>. Pour le Hive metastore, il doit être de la forme <schema>.<table>.

  • L'exécution d'une mise à jour de refresh complet ne nettoie pas les données de résultats calculées précédemment dans les récepteurs. Cela signifie que toutes les données retraitées sont ajoutées au récepteur, et les données existantes ne sont pas modifiées.

  • Les attentes de pipeline ne sont pas prises en charge.

  • Le contrôle de sortie Serverless prend en charge uniquement les connecteurs de destination Kafka et Delta Lake. Voir Qu'est-ce que le contrôle de sortie Serverless ?.

Ressources supplémentaires