Configurer le mode temps réel
Cette page décrit les prérequis et la configuration nécessaires pour exécuter des requêtes en mode temps réel dans Structured Streaming. Pour un didacticiel pas à pas, consultez Didacticiel : Exécuter une charge de travail de streaming en temps réel. Pour des informations conceptuelles sur le mode temps réel, consultez Mode temps réel dans Structured Streaming.
Prérequis
Pour utiliser le mode temps réel, vous devez configurer votre compute afin de satisfaire aux exigences suivantes :
- Utilisez le compute classique. Les modes d'accès dédié et standard sont pris en charge. Le mode d'accès standard est pris en charge pour Python uniquement. Les LakeFlow Pipelines et les clusters Serverless ne sont pas pris en charge.
- Utiliser Databricks Runtime 16.4 LTS et versions supérieures.
- Désactiver la mise à l'échelle automatique.
- Désactiver Photon.
- Définir
spark.databricks.streaming.realTimeMode.enabledsurtrue. - Désactivez les instances ponctuelles pour éviter les interruptions.
Pour les charges de travail sensibles à la latence avec des UDF, Databricks vous recommande d'utiliser le mode d'accès dédié. Voir Fonctions de table.
Pour obtenir des instructions sur la création et la configuration du compute classique, consultez Référence de configuration du compute.
Jointures Stream à Stream
Les jointures internes Stream à Stream nécessitent une configuration supplémentaire pour le mode temps réel. Les jointures externes ne sont pas prises en charge. Consultez Stream to stream join.
Pour exécuter une jointure stream to stream en mode temps réel avec plusieurs autres streams sur le même cluster, vous devez utiliser Databricks Runtime 18 et versions supérieures.
Dans Databricks Runtime 18.2 et versions antérieures, Structured Streaming ne prend pas en charge les configurations suivantes pour les autres modes de traitement, notamment processingTime et availableNow.
Pour activer les jointures Stream vers Stream en mode temps réel, définissez les configurations Spark suivantes :
- Python
- Scala
- SQL
spark.conf.set("spark.databricks.streaming.realTimeMode.streamStreamJoin.enabled", "true")
spark.conf.set("spark.sql.streaming.join.stateFormatVersion", "4")
spark.conf.set("spark.sql.streaming.join.stateFormatV4.enabled", "true")
spark.conf.set("spark.sql.streaming.stateStore.rocksdb.mergeOperatorVersion", "2")
spark.conf.set("spark.sql.streaming.realTimeMode.controlMessage.enabled", "true")
spark.conf.set("spark.databricks.streaming.realTimeMode.streamStreamJoin.enabled", "true")
spark.conf.set("spark.sql.streaming.join.stateFormatVersion", "4")
spark.conf.set("spark.sql.streaming.join.stateFormatV4.enabled", "true")
spark.conf.set("spark.sql.streaming.stateStore.rocksdb.mergeOperatorVersion", "2")
spark.conf.set("spark.sql.streaming.realTimeMode.controlMessage.enabled", "true")
SET spark.databricks.streaming.realTimeMode.streamStreamJoin.enabled = true;
SET spark.sql.streaming.join.stateFormatVersion = 4;
SET spark.sql.streaming.join.stateFormatV4.enabled = true;
SET spark.sql.streaming.stateStore.rocksdb.mergeOperatorVersion = 2;
SET spark.sql.streaming.realTimeMode.controlMessage.enabled = true;
Configuration de la query
Pour exécuter une query en mode temps réel, vous devez activer le Trigger temps réel. Les Trigger temps réel ne sont pris en charge qu’en mode mise à jour.
- Python
- Scala
query = (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", broker_address)
.option("subscribe", input_topic)
.load()
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", broker_address)
.option("topic", output_topic)
.option("checkpointLocation", checkpoint_location)
.outputMode("update")
# In PySpark, the realTime trigger requires specifying the interval.
.trigger(realTime="5 minutes")
.start()
)
import org.apache.spark.sql.execution.streaming.RealTimeTrigger
val readStream = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", brokerAddress)
.option("subscribe", inputTopic).load()
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", brokerAddress)
.option("topic", outputTopic)
.option("checkpointLocation", checkpointLocation)
.outputMode("update")
.trigger(RealTimeTrigger.apply())
// RealTimeTrigger can also accept an argument specifying the checkpoint interval.
// For example, this code indicates a checkpoint interval of 5 minutes:
// .trigger(RealTimeTrigger.apply("5 minutes"))
.start()
Dimensionnement du compute
Vous pouvez exécuter un Job en temps réel par ressource de compute si le compute dispose de suffisamment d'emplacements de tâches.
Pour fonctionner en mode faible latence, le nombre total d'emplacements de tâche disponibles doit être supérieur ou égal au nombre de tâches à travers toutes les étapes de query.
Exemples de calcul d'emplacement
Type de pipeline | Configuration | Emplacements requis |
|---|---|---|
Sans état à un seul étage (source Kafka + sink) |
| 8 emplacements |
À deux étapes et avec état (source Kafka + brassage) |
| 28 slots (8 + 20) |
À trois étapes (source Kafka + shuffle + repartition) |
| 48 emplacements (8 + 20 + 20) |
Si vous ne définissez pas maxPartitions, utilisez le nombre de partitions dans le sujet Kafka.