Aller au contenu principal

query streaming data

Vous pouvez utiliser Databricks pour query les sources de données en streaming à l'aide de Structured Streaming. Databricks prend en charge de manière exhaustive les charges de travail en streaming en Python et Scala, et prend en charge la plupart des fonctionnalités de Structured Streaming avec SQL.

Les exemples suivants montrent comment utiliser un récepteur de mémoire pour l'inspection manuelle des données de streaming pendant le développement interactif dans les Notebooks. En raison des limites de sortie des lignes dans l'interface utilisateur du Notebook, vous ne pourrez peut-être pas observer toutes les données lues par les requêtes de streaming. Dans les charges de travail de production, vous ne devriez Trigger des queries de streaming qu'en les écrivant dans une table cible ou un système externe.

remarque

La prise en charge SQL des requêtes interactives sur les données de streaming est limitée aux notebooks s'exécutant sur des compute multifonctions. Vous pouvez également utiliser SQL lorsque vous déclarez des tables de streaming dans Databricks SQL ou des pipelines LakeFlow Pipelines. Consultez Tables de streaming et Spark Declarative Pipelines.

query data from streaming systems

Databricks fournit des lecteurs de données en streaming pour les systèmes de streaming suivants :

  • Kafka
  • Kinesis
  • PubSub
  • Pulsar

Vous devez fournir les détails de configuration lorsque vous initialisez des queries sur ces systèmes, qui varient en fonction de votre environnement configuré et du système à partir duquel vous choisissez de lire. Consultez les connecteurs standard de LakeFlow Connect.

Les charges de travail courantes impliquant des systèmes de streaming incluent l'ingestion de données vers le lakehouse et le traitement de Stream pour envoyer des données vers des systèmes externes. Pour en savoir plus sur les charges de travail de streaming, consultez les concepts de Structured Streaming.

Les exemples suivants démontrent une lecture streaming interactive à partir de Kafka :

Python
display(spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "latest")
.load()
)

query une table en lecture streaming

Databricks crée toutes les tables en utilisant Delta Lake par default. Lorsque vous exécutez une query de streaming sur une table Delta, la query récupère automatiquement les nouveaux enregistrements lorsqu’une version de la table est validée. Par default, les queries streaming s'attendent à ce que les tables sources ne contiennent que des enregistrements ajoutés. Si vous devez travailler avec des données streaming qui contiennent des mises à jour et des suppressions, Databricks recommande d'utiliser les LakeFlow Pipelines et AUTO CDC ... INTO. Voir les APIs AUTO CDC : Simplifiez la capture des changements de données avec les pipelines.

Les exemples suivants montrent comment effectuer une lecture en streaming interactive à partir d'une table :

Python
display(spark.readStream.table("table_name"))

Interroger les données du stockage d'objets cloud avec Auto Loader

Vous pouvez Stream des données à partir du stockage d’objets cloud à l’aide d’Auto Loader, le connecteur de données cloud de Databricks. Vous pouvez utiliser le connecteur avec des fichiers stockés dans les volumes Unity Catalog ou d’autres emplacements de stockage d’objets cloud. Databricks recommande d’utiliser des volumes pour gérer l’accès aux données dans le stockage d’objets cloud. Voir Connecter à des sources de données et services externes.

Databricks optimise ce connecteur pour l'ingestion en streaming de données dans le stockage d'objets cloud qui sont stockées dans des formats structurés, semi-structurés et non structurés populaires. Databricks recommande de stocker les données ingérées dans un format quasi-brut pour maximiser le throughput et minimiser la perte de données potentielle due à des enregistrements corrompus ou des modifications de schéma.

Pour plus de recommandations sur l'ingestion de données à partir du stockage d'objets cloud, consultez Connecteurs standard dans Lakeflow Connect.

Les exemples suivants illustrent une lecture streaming interactive à partir d'un répertoire de fichiers JSON dans un volume :

Python
display(spark.readStream.format("cloudFiles").option("cloudFiles.format", "json").load("/Volumes/catalog/schema/volumes/path/to/files"))