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.
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
- SQL
display(spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "latest")
.load()
)
SELECT * FROM STREAM read_kafka(
bootstrapServers => '<server:ip>',
subscribe => '<topic>',
startingOffsets => 'latest'
);
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
- SQL
display(spark.readStream.table("table_name"))
SELECT * FROM STREAM 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
- SQL
display(spark.readStream.format("cloudFiles").option("cloudFiles.format", "json").load("/Volumes/catalog/schema/volumes/path/to/files"))
SELECT * FROM STREAM read_files('/Volumes/catalog/schema/volumes/path/to/files', format => 'json')