Remarque
L’accès à cette page requiert une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page requiert une autorisation. Vous pouvez essayer de modifier des répertoires.
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 obtenir un didacticiel pas à pas, consultez Tutoriel : Exécuter une charge de travail de streaming en temps réel. Pour plus d’informations conceptuelles sur le mode en temps réel, consultez le mode temps réel dans Structured Streaming.
Conditions préalables
Pour utiliser le mode en temps réel, vous devez configurer votre calcul pour répondre aux exigences suivantes :
- Utilisez le calcul classique. Les modes d’accès dédiés et standard sont pris en charge. Le mode d’accès standard est pris en charge uniquement pour Python. Les pipelines Lakeflow et les clusters serverless ne sont pas pris en charge.
- Utilisez Databricks Runtime 16.4 LTS et versions ultérieures.
- Désactivez la mise à l’échelle automatique.
- Désactivez Photon.
- Affectez la valeur
spark.databricks.streaming.realTimeMode.enabledàtrue. - Désactivez les instances spot pour éviter les interruptions.
Pour les charges de travail sensibles à la latence avec des fonctions UDF, Databricks recommande d’utiliser le mode d'accès dédié. Consultez les fonctions de la table.
Pour obtenir des instructions sur la création et la configuration du calcul classique, consultez la référence de configuration de calcul.
Jointures flux-flux
Les jointures internes entre flux nécessitent une configuration supplémentaire en mode temps réel. Les jointures externes ne sont pas prises en charge. Consultez la section Jointures flux-flux.
Important
Pour exécuter une jointure flux à flux en mode temps réel avec plusieurs autres flux sur le même cluster, vous devez utiliser Databricks Runtime 18 LTS ou version ultérieure.
Dans Databricks Runtime 18.2 et ci-dessous, Structured Streaming ne prend pas en charge les configurations suivantes pour d’autres modes de traitement, notamment processingTime et availableNow.
Pour activer les jointures flux à flux en mode temps réel, définissez les configurations Spark suivantes :
Python
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")
Scala
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")
SQL
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 des requêtes
Pour exécuter une requête en mode temps réel, vous devez activer le déclencheur en temps réel. Les déclencheurs en temps réel sont pris en charge uniquement en mode mise à jour.
Python
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()
)
Scala
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 des ressources de calcul
Vous pouvez exécuter un travail en temps réel par ressource de calcul si le calcul a suffisamment d’emplacements de tâches.
Pour s’exécuter en mode faible latence, le nombre total d’emplacements de tâches disponibles doit être supérieur ou égal au nombre de tâches dans toutes les phases de requête.
Exemples de calcul d'emplacement
| Type de pipeline | Paramétrage | Emplacements requis |
|---|---|---|
| Étape unique sans état (source Kafka + récepteur) |
maxPartitions = 8 |
8 emplacements |
| Avec état à deux étapes (source Kafka + shuffle) |
maxPartitions = 8, partitions en shuffle = 20 |
28 emplacements (8 + 20) |
| Trois étapes (source Kafka + shuffle + répartition) |
maxPartitions = 8, deux phases en shuffle de 20 chacune |
48 emplacements (8 + 20 + 20) |
Si vous ne définissez maxPartitionspas, utilisez le nombre de partitions dans la rubrique Kafka.