Aller au contenu principal

Streaming sur compute serverless

Cette page décrit comment choisir la bonne configuration pour les charges de travail de streaming Serverless sur Databricks, y compris les pipelines continus, l'ingestion incrémentielle et les connecteurs gérés. Le choix de la bonne configuration dépend de la source, de la forme et des besoins en latence du Stream.

Qu'est-ce qui constitue une charge de travail en streaming

Une charge de travail de streaming lit des données non bornées à partir d'une source (tels que le stockage d'objets cloud, un bus de messages ou un flux de changements) et écrit de manière incrémentielle dans un récepteur. Databricks prend en charge deux modèles de charges de travail de streaming :

  • **Continu** : Un pipeline qui s'exécute sans interruption et traite les nouvelles données dès leur arrivée. La latence est mesurée en secondes.
  • Incrémentiel (également appelé déclenché) : Un pipeline qui s'exécute selon un calendrier ou un Trigger, traite toutes les données arrivées depuis la dernière exécution, puis s'arrête. La latence est mesurée en minutes.

Certaines charges de travail semblent être des pipelines de streaming, mais ne sont pas techniquement des pipelines. Les exemples incluent un service qui maintient un websocket ouvert pour écouter les événements, une application de chat qui maintient une connexion persistante par utilisateur, ou un récepteur de webhook qui gère les requêtes HTTP entrantes. Ce sont des applications, pas des pipelines de streaming. Pour la bonne option serverless pour ces charges de travail, consultez Charges de travail qui ne sont pas des pipelines de streaming.

Choisissez la bonne configuration de streaming

Ce tableau mappe les cas d'utilisation aux configurations serverless qui leur conviennent le mieux. Les sections qui suivent sur cette page fournissent plus de détails sur ces recommandations.

Cas d'usage

Configuration recommandée

Pourquoi

Streaming ETL continu à faible latence ou Transformations

LakeFlow Pipelines en mode continu

Le mode continu est conçu pour les Stream toujours actifs. Le pipelining de Stream exécute des micro-lots simultanément, améliorant le throughput et la latence. L'état géré maintient la récupération automatique.

Ingestion incrémentielle à partir du stockage cloud

Utilisez Auto Loader dans les LakeFlow Pipelines (pour une faible latence) ou dans un Job serverless avec Trigger.AvailableNow() (si une latence plus faible est acceptable).

Auto Loader suit efficacement les nouveaux fichiers. Trigger.AvailableNow() traite le backlog, puis se termine, ce qui correspond à une cadence planifiée ou à la demande.

Ingestion gérée à partir de sources SaaS ou de CDC de base de données

Connecteurs standard dans Lakeflow Connect

Connecteurs entièrement managés avec des pipelines d'ingestion serverless. Aucun code n'est requis pour les sources prises en charge.

Streaming SQL sur les tables Delta

Tables de streaming

Traitement incrémentiel natif SQL pour les sources orientées ajout, avec des pipelines gérés et du refresh.

Traitement périodique par micro-batch dans un Notebook ou un job

Job serverless avec Trigger.AvailableNow()

Rentable lorsque la fraîcheur à la minute est suffisante. Le Compute Serverless start rapidement et s'arrête lorsque le batch se termine.

Cas d'usage

Configuration recommandée

Pourquoi

Streaming ETL continu à faible latence ou Transformations

LakeFlow Pipelines en mode continu

Le mode continu est conçu pour les Stream toujours actifs. Le pipelining de Stream exécute des micro-lots simultanément, améliorant le throughput et la latence. L'état géré maintient la récupération automatique.

Ingestion incrémentielle à partir du stockage cloud

Utilisez Auto Loader dans les LakeFlow Pipelines (pour une faible latence) ou dans un Job serverless avec Trigger.AvailableNow() (si une latence plus faible est acceptable).

Auto Loader suit efficacement les nouveaux fichiers. Trigger.AvailableNow() traite le backlog, puis se termine, ce qui correspond à une cadence planifiée ou à la demande.

Ingestion gérée à partir de sources SaaS ou de CDC de base de données

Connecteurs standard dans Lakeflow Connect

Connecteurs entièrement managés avec des pipelines d'ingestion serverless. Aucun code n'est requis pour les sources prises en charge.

Streaming SQL sur les tables Delta

Tables de streaming

Traitement incrémentiel natif SQL pour les sources orientées ajout, avec des pipelines gérés et du refresh.

Traitement périodique par micro-batch dans un Notebook ou un job

Job serverless avec Trigger.AvailableNow()

Rentable lorsque la fraîcheur à la minute est suffisante. Le Compute Serverless start rapidement et s'arrête lorsque le batch se termine.

Streaming continu

Pour le streaming continu sur compute Serverless, utilisez les LakeFlow Pipelines en mode continu. Le pipeline reste en cours d'exécution, traite les enregistrements au fur et à mesure de leur arrivée et récupère automatiquement les échecs.

Pour configurer un Stream continu :

astuce

Le Stream pipelining est activé par default dans les pipelines Lakeflow serverless. Les micro-lots s'exécutent simultanément plutôt que séquentiellement, ce qui améliore le throughput des streams à forte ingestion.

Les triggers Structured Streaming basés sur le temps, tels que Trigger.ProcessingTime(interval) et Trigger.Continuous(interval), ne sont pas disponibles dans les Notebooks ou les Jobs serverless. Utilisez les Lakeflow pipelines en mode continu pour le modèle de fonctionnement permanent. Consultez les limitations du streaming. Trigger.Once() est pris en charge mais obsolète — migrez les requêtes existantes vers Trigger.AvailableNow().

Streaming incrémental et déclenché

Pour le streaming incrémentiel, exécutez Structured Streaming avec Trigger.AvailableNow() dans un Job Serverless. Chaque exécution traite toutes les données arrivées depuis le dernier point de contrôle, puis se termine.

Pour configurer un job serverless avec streaming incrémentiel :

L'exemple suivant lit de nouveaux fichiers à partir du stockage cloud (source_path) avec Auto Loader, traite toutes les données disponibles au moment de l'exécution et écrit dans une table Delta :

Python
(spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.maxFilesPerTrigger", 1000)
.load(source_path)
.writeStream
.trigger(availableNow=True)
.option("checkpointLocation", checkpoint_path)
.toTable("catalog.schema.target_table"))

Un Job Trigger.AvailableNow() planifié est le modèle de streaming le plus rentable sur le compute Serverless lorsque la latence d'une minute est acceptable. Le compute start en quelques secondes, exécute le batch et s'arrête.

Ingestion gérée

Si la source est une application SaaS ou une base de données opérationnelle, utilisez Lakeflow Connect plutôt que d'écrire du code Structured Streaming. Lakeflow Connect exécute des pipelines d'ingestion serverless pour des connecteurs tels que Salesforce, Workday, SQL Server CDC et PostgreSQL CDC. Consultez les connecteurs gérés dans Lakeflow Connect.

Ce chemin est la bonne réponse lorsque :

  • Un connecteur existe pour votre source.
  • Vous souhaitez un pipeline géré plutôt qu'un code personnalisé.
  • Vous avez besoin d'une évolution des schémas, d'une lignée et d'un monitoring prêts à l'emploi.

Traitement de données incrémentiel géré par SQL

Pour les équipes axées sur SQL, utilisez les tables de streaming pour les charges de travail de streaming natives SQL. Vous pouvez définir des tables de streaming au sein des LakeFlow Pipelines ou en tant que tables de streaming autonomes.

Pour les tables de streaming autonomes créées avec l'instruction SQL CREATE OR REFRESH STREAMING TABLE, le refresh initial des données et le peuplement commencent immédiatement. Un pipeline serverless dédié est automatiquement créé et géré par le système pour chaque table de streaming.

Si vous avez besoin de résultats de query avec des sémantiques de batch et un refresh géré, utilisez plutôt des vues matérialisées. Consultez les vues matérialisées.

Charges de travail qui ne sont pas des pipelines de streaming

Une charge de travail qui doit maintenir une connexion persistante, écouter sur un port ou répondre aux requêtes HTTP entrantes n’est pas un pipeline de streaming ; il s’agit d’une application. N’exécutez pas ces charges de travail sur un Job serverless. Les bonnes options Databricks sont :

  • Services de longue durée nécessitant une connexion persistante ou un Endpoint HTTP : Développez le service avec Databricks Apps. Databricks Apps est la plateforme Serverless pour héberger des applications personnalisées sur Databricks, y compris FastAPI, Flask, Streamlit, Dash, Gradio, Node.js et les applications Shiny. See Databricks Apps.
  • Webhooks entrants ou écouteurs d'événements : Exposez un endpoint HTTP sur Databricks Apps ou mettez fin au webhook dans un service externe et écrivez les événements vers le stockage cloud ou un bus de messages, puis récupérez-les avec un pipeline de streaming serverless.
  • Échange de jetons ou d'informations d'identification personnalisés : utilisez des Service Principal avec OAuth, ou appelez les Databricks REST APIs à partir d'une application. Les pipelines de streaming ne conservent pas les sessions par utilisateur ni l'état des jetons personnalisés.

Si vous évaluez si votre charge de travail s'adapte à un pipeline de streaming, demandez :

  • La charge de travail lit-elle à partir d'une source de données illimitée et écrit-elle vers un récepteur ? Si oui, il s'agit d'un pipeline de streaming.
  • La charge de travail doit-elle maintenir une connexion ouverte à un client ? Si oui, il s'agit d'une application ; utilisez Databricks Apps.

Limitations

Le compute serverless impose les contraintes de streaming suivantes. Aucun d'entre eux n'empêche les charges de travail ci-dessus lorsqu'ils sont associés au bon produit.

  • Les Trigger Structured Streaming basés sur le temps (Trigger.ProcessingTime(interval) et Trigger.Continuous(interval)) ne sont pas pris en charge dans les Notebooks ou Jobs Serverless. Utilisez les Lakeflow pipelines en mode continu pour les Streams toujours actifs, ou Trigger.AvailableNow() pour les exécutions déclenchées. Consultez les limitations du streaming.
  • Les **streaming queries** sans **Trigger** explicite échouent avec INFINITE_STREAMING_TRIGGER_NOT_SUPPORTED. Apache Spark utilise default Trigger.ProcessingTime("0 seconds"), ce qui n'est pas pris en charge sur le compute serverless. Définissez toujours Trigger.AvailableNow() sur chaque **streaming query**, ou utilisez les **LakeFlow Pipelines** en mode continu.
  • Toutes les limitations concernant le streaming en mode d'accès standard s'appliquent également au compute Serverless. Consultez les limitations de streaming.

Étapes suivantes