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.
- L'API
sinkest uniquement disponible pour Python. - Vous pouvez créer un récepteur personnalisé avec l'API ForEachBatch. Voir Utiliser ForEachBatch pour écrire dans des stockages de données arbitraires dans des pipelines.
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 :
- Créer le puits.
- 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 :
- Delta sinks
- Kafka and Azure Event Hubs sinks
- Python custom data sources
Pour créer un récepteur Delta par chemin de fichier :
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é :
dp.create_sink(
name = "delta_sink",
format = "delta",
options = { "tableName": "catalog_name.schema_name.table_name" }
)
Ce code fonctionne pour les destinations Apache Kafka et Azure Event Hub.
credential_name = "<service-credential>"
eh_namespace_name = "dp-eventhub"
bootstrap_servers = f"{eh_namespace_name}.servicebus.windows.net:9093"
topic_name = "dp-sink"
dp.create_sink(
name = "eh_sink",
format = "kafka",
options = {
"databricks.serviceCredential": credential_name,
"kafka.bootstrap.servers": bootstrap_servers,
"topic": topic_name
}
)
Le credential_name est une référence à un identifiant de service Unity Catalog. Pour plus d'informations, consultez Utiliser les informations d'identification de service Unity Catalog pour se connecter à des services cloud externes.
En supposant que vous ayez une source de données Python personnalisée enregistrée sous le nom my_custom_datasource, le code suivant peut alors écrire dans cette source de données.
from pyspark import pipelines as dp
# Assume `my_custom_datasource` is a custom Python streaming
# data source that writes data to your system.
# Create Lakeflow pipelines sink using my_custom_datasource
dp.create_sink(
name="custom_sink",
format="my_custom_datasource",
options={
<options-needed-for-custom-datasource>
}
)
# Create append flow to send data to RequestBin
@dp.append_flow(name="flow_to_custom_sink", target="custom_sink")
def flow_to_custom_sink():
return read_stream("my_source_data")
Pour plus de détails sur la création de sources de données personnalisées en Python, consultez les sources de données personnalisées PySpark.
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
deltaet 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
kafkaet 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
kafkaet 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.
- Delta sink
- Kafka and Azure Event Hubs sinks
@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")
)
@dp.append_flow(name = "kafka_sink_flow", target = "eh_sink")
def kafka_sink_flow():
return (
spark.readStream.table("spark_referrers")
.selectExpr("cast(current_page_id as string) as key", "to_json(struct(referrer, current_page_title, click_count)) AS value")
)
Le paramètre value est obligatoire pour une cible Azure Event Hub. Des paramètres supplémentaires tels que key, partition, headers et topic sont facultatifs.
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_flowetupdate_flowpeuvent être utilisés pour écrire dans des récepteurs. Les autres flux, tels quecreate_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 ?.