create_sink
La fonction create_sink() écrit dans un service de streaming d'événements tel que Apache Kafka ou Azure Event Hub ou dans une table Delta à partir d'un pipeline déclaratif. Après avoir créé un récepteur avec la fonction create_sink(), vous utilisez le récepteur dans un flux d'ajout ou un flux de mise à jour pour écrire des données dans le récepteur. Les flux d'ajout et de mise à jour sont les seuls types de flux pris en charge avec la fonction create_sink(). D'autres types de flux, tels que create_auto_cdc_flow, ne sont pas pris en charge. Pour plus de détails sur les autres types de récepteurs dans les LakeFlow Pipelines, consultez Récepteurs dans les LakeFlow Pipelines.
Le récepteur Delta prend en charge les tables externes et gérées Unity Catalog et les tables gérées du Hive metastore. Les noms de table doivent être entièrement qualifiés. Par exemple, les tables Unity Catalog doivent utiliser un identifiant à trois niveaux : <catalog>.<schema>.<table>. Les tables du Hive metastore doivent utiliser <schema>.<table>.
- L'exécution d'une mise à jour de refresh complète ne supprime pas les données des cibles. Toutes les données retraitées sont ajoutées à la destination, et les données existantes ne sont pas modifiées.
- Les attentes ne sont pas prises en charge avec l'API
sink.
Syntaxe
from pyspark import pipelines as dp
dp.create_sink(name=<sink_name>, format=<format>, options=<options>)
parameter
parameter | Type | Description |
|---|---|---|
|
| Obligatoire. Une chaîne qui identifie le récepteur et est utilisée pour le référencer et le gérer. Les noms du récepteur doivent être uniques pour le pipeline, y compris pour tous les fichiers de code source qui font partie du pipeline. |
|
| Obligatoire. Une chaîne qui définit le format de sortie, soit |
|
| Une liste d'options de sink, formatée comme
|
Exemples
from pyspark import pipelines as dp
# Create a Kafka sink
dp.create_sink(
"my_kafka_sink",
"kafka",
{
"kafka.bootstrap.servers": "host:port",
"topic": "my_topic"
}
)
# Create an external Delta table sink with a file path
dp.create_sink(
"my_delta_sink",
"delta",
{ "path": "/path/to/my/delta/table" }
)
# Create a Delta table sink using a table name
dp.create_sink(
"my_delta_sink",
"delta",
{ "tableName": "my_catalog.my_schema.my_table" }
)