Exécutez votre première charge de travail Structured Streaming
Cet article fournit des exemples de code et une explication des concepts de base nécessaires pour exécuter vos premières queries Structured Streaming sur Databricks. Vous pouvez utiliser Structured Streaming pour les charges de travail de traitement quasi temps réel et incrémentiel.
Structured Streaming est l'une des nombreuses technologies qui alimentent les tables de streaming dans les Lakeflow pipelines. Databricks recommande d'utiliser les Lakeflow Pipelines pour toutes les nouvelles charges de travail ETL, d'ingestion et de Structured Streaming. See Spark Declarative Pipelines.
Alors que les Lakeflow pipelines fournissent une syntaxe légèrement modifiée pour déclarer les tables de streaming, la syntaxe générale pour la configuration des lectures et transformations de streaming s'applique à tous les cas d'utilisation de streaming sur Databricks. Les Lakeflow Pipelines simplifient également le streaming en gérant les informations d'état, les métadonnées et de nombreuses configurations.
Utilisez Auto Loader pour lire les données en streaming depuis le stockage d'objets
L'exemple suivant montre le chargement de données JSON avec Auto Loader, qui utilise cloudFiles pour désigner le format et les options. L'option schemaLocation permet l'inférence et l'évolution des schémas. Collez le code suivant dans une cellule de notebook Databricks et exécutez la cellule pour créer un DataFrame de streaming nommé raw_df:
file_path = "/databricks-datasets/structured-streaming/events"
checkpoint_path = "/tmp/ss-tutorial/_checkpoint"
raw_df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", checkpoint_path)
.load(file_path)
)
Comme les autres opérations de lecture sur Databricks, la configuration d'une lecture en streaming ne charge pas réellement les données. Vous devez Trigger une action sur les données avant que le Stream ne commence.
Appeler display() sur une DataFrame de streaming démarre un Job de streaming. Pour la plupart des cas d'utilisation de Structured Streaming, l'action qui déclenche un Stream devrait écrire des données dans un récepteur. Consultez les considérations de production pour Structured Streaming.
Effectuer une transformation en streaming
Structured Streaming prend en charge la plupart des Transformations disponibles dans Databricks et Spark SQL. Vous pouvez même charger des modèles MLflow en tant que fonctions UDF et faire des prédictions en streaming sous forme de transformation.
L'exemple de code suivant effectue une simple transformation pour enrichir les données JSON ingérées avec des informations supplémentaires à l'aide des fonctions Spark SQL :
from pyspark.sql.functions import col, current_timestamp
transformed_df = (raw_df.select(
"*",
col("_metadata.file_path").alias("source_file"),
current_timestamp().alias("processing_time")
)
)
Le transformed_df résultant contient des instructions de query pour charger et transformer chaque enregistrement tel qu'il arrive dans la source de données.
Structured Streaming traite les sources de données comme des datasets illimités ou infinis. Ainsi, certaines transformations ne sont pas prises en charge dans les charges de travail de Structured Streaming, car elles nécessiteraient le tri d'un nombre infini d'éléments.
La plupart des agrégations et de nombreuses jointures nécessitent la gestion des informations d'état avec des filigranes, des fenêtres et un mode de sortie. Consultez Appliquer des filigranes pour contrôler les seuils de traitement des données.
Effectuer une écriture par batch incrémentielle dans Delta Lake
L’exemple suivant écrit dans Delta Lake à l’aide d’un chemin de fichier et d’un point de contrôle spécifiés.
Assurez-vous toujours de spécifier un emplacement de point de contrôle unique pour chaque writer de streaming que vous configurez. Le point de contrôle fournit l'identité unique de votre stream, en assurant le suivi de tous les enregistrements traités et des informations d'état associées à votre query streaming.
Le paramètre availableNow du Trigger indique à Structured Streaming de traiter tous les enregistrements précédemment non traités du dataset source, puis de s'arrêter, afin que vous puissiez exécuter en toute sécurité le code suivant sans vous soucier de laisser un Stream en cours d'exécution :
target_path = "/tmp/ss-tutorial/"
checkpoint_path = "/tmp/ss-tutorial/_checkpoint"
transformed_df.writeStream
.trigger(availableNow=True)
.option("checkpointLocation", checkpoint_path)
.option("path", target_path)
.start()
Dans cet exemple, aucun nouvel enregistrement n’arrive dans notre source de données, de sorte que l’exécution répétée de ce code n’ingère pas de nouveaux enregistrements.
L'exécution de Structured Streaming peut empêcher l'arrêt automatique des ressources de compute. Pour éviter des coûts inattendus, assurez-vous de terminer les requêtes de streaming.
Lire les données depuis Delta Lake, transformer et écrire dans Delta Lake
Delta Lake prend en charge de manière approfondie l'utilisation de Structured Streaming en tant que source et récepteur. Consultez Lectures et écritures en streaming de tables Delta Lake.
L'exemple suivant présente la syntaxe d'exemple pour charger de manière incrémentielle tous les nouveaux enregistrements d'une table Delta Lake, les joindre avec un instantané d'une autre table Delta Lake et les écrire dans une table Delta Lake :
(spark.readStream
.table("<table-name1>")
.join(spark.read.table("<table-name2>"), on="<id>", how="left")
.writeStream
.trigger(availableNow=True)
.option("checkpointLocation", "<checkpoint-path>")
.toTable("<table-name3>")
)
Vous devez disposer des autorisations appropriées configurées pour lire les tables sources et écrire dans les tables cibles ainsi que l'emplacement de point de contrôle spécifié. Remplissez tous les parameter désignés par des crochets obliques (<>) à l'aide des valeurs pertinentes pour vos sources et destinations de données.
Les LakeFlow Pipelines fournissent une syntaxe entièrement déclarative pour créer des pipelines Delta Lake et gérer automatiquement des propriétés comme les triggers et les checkpoints. See Spark Declarative Pipelines.
Lire les données de Kafka, transformer et écrire dans Kafka
Apache Kafka et d'autres bus de messagerie offrent l'une des latences les plus faibles disponibles pour les grands datasets. Vous pouvez utiliser Databricks pour appliquer des transformations aux données ingérées depuis Kafka, puis réécrire les données dans Kafka.
L'écriture de données vers un stockage d'objets cloud ajoute une surcharge de latence supplémentaire. Si vous souhaitez stocker des données provenant d'un bus de messagerie dans Delta Lake, mais que vous avez besoin de la latence la plus faible possible pour les workloads de streaming, Databricks vous recommande de configurer des jobs de streaming distincts pour ingérer les données dans le lakehouse et d'appliquer des transformations quasi en temps réel pour les cibles de bus de messagerie en aval.
Le code d'exemple suivant démontre un modèle simple pour enrichir des données de Kafka en les joignant avec des données d'une table Delta Lake, puis en les réécrivant dans Kafka :
(spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "latest")
.load()
.join(spark.read.table("<table-name>"), on="<id>", how="left")
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("topic", "<topic>")
.option("checkpointLocation", "<checkpoint-path>")
.start()
)
Vous devez disposer des autorisations appropriées configurées pour l'accès à votre service Kafka. Remplissez tous les paramètres indiqués par des crochets d’angle (<>) en utilisant les valeurs pertinentes pour vos sources de données et vos puits de données. Consultez Connectez-vous à Apache Kafka.