Sinks dans les LakeFlow Pipelines
Par défaut, les flux du pipeline écrivent les résultats dans des tables Delta gérées par Unity Catalog, généralement des tables de streaming ou des vues matérialisées. Les cibles de sortie sont une cible alternative qui vous permet d'écrire des données transformées vers des destinations situées en dehors du stockage géré par Databricks, telles que des services de streaming d'événements ou des magasins de données personnalisés.
Les sinks sont utilisés avec les flux d'ajout. Vous définissez un puits à l'aide de l'une des APIs de puits, puis vous y faites référence en tant que target dans votre définition append_flow.
Quand utiliser des récepteurs
Databricks recommande d'utiliser des récepteurs lorsque vous en avez besoin :
- Créez des cas d'utilisation opérationnels à faible latence, tels que la détection de la fraude, l'analytique en temps réel ou les recommandations client, où les données doivent circuler vers un bus de messages plutôt que vers un stockage cloud. Pour les charges de travail qui nécessitent une latence de l'ordre de la milliseconde, consultez Utiliser le mode temps réel dans les pipelines Lakeflow.
- Écrire les données transformées dans des tables gérées par une instance Delta externe, y compris les tables gérées et externes de Unity Catalog.
- Effectuez l'ETL inverse vers des systèmes externes, par exemple en réécrivant les données traitées dans les rubriques Apache Kafka pour une consommation en dehors de Databricks.
- Écrivez dans un format non pris en charge en mode natif par Databricks, à l’aide de sources de données Python personnalisées.
Types de récepteur
Les pipelines prennent en charge les types de destination suivants :
Type de puits | Description |
|---|---|
Récepteurs de table Delta | Écrire dans des tables Delta gérées par Unity Catalog ou externes. Spécifiez un chemin de fichier ou un nom de table entièrement qualifié. |
Puits Apache Kafka | Écrivez dans les rubriques Apache Kafka à l'aide du connecteur Kafka inclus dans l'environnement d'exécution du pipeline. |
Destinations Azure Event Hub | Écrire dans Azure Event Hubs en utilisant l'interface Kafka. Utilise les mêmes options que les récepteurs Kafka. |
Puits Python personnalisés | Écrivez dans n’importe quelle source de données à l’aide d’une source de données personnalisée Python enregistrée avec |
Puits ForEachBatch | Appliquez une logique Python personnalisée à chaque micro-batch de données en streaming. Utilisez lorsque vous devez écrire vers plusieurs destinations, effectuer des upserts ou utiliser des cibles qui ne prennent pas en charge les écritures en streaming en mode natif. |
APIs Sink
Les pipelines fournissent deux API pour créer des récepteurs :
create_sink(): Crée un récepteur nommé d'un type pris en charge (Delta, Kafka, AEH ou source de données Python personnalisée). Uniquement disponible dans Python. Voir Utiliser des puits dans des pipelines.foreach_batch_sink(): Décore une fonction Python qui s’exécute pour chaque micro-batch de données en streaming. Offre une flexibilité maximale pour une logique d'écriture personnalisée. Voir Utiliser ForEachBatch pour écrire dans des dépôts de données arbitraires dans les pipelines.
Les deux types de sink sont référencés comme le target d'un append_flow.
Limitations
- Les récepteurs ne sont disponibles qu'en Python. 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 les
append_flowpeuvent écrire vers des récepteurs ; lescreate_auto_cdc_flowet autres types de flux ne sont pas pris en charge. - Les attentes de pipeline ne sont pas prises en charge pour les puits.
- L'exécution d'un refresh complet ne nettoie pas les données précédemment écrites dans les récepteurs.