Aller au contenu principal

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.enabled sur true.
  • 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.

important

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
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")

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
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()
)

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)

maxPartitions = 8

8 emplacements

À deux étapes et avec état (source Kafka + brassage)

maxPartitions = 8, partitions de shuffle = 20

28 slots (8 + 20)

À trois étapes (source Kafka + shuffle + repartition)

maxPartitions = 8, deux étapes de brassage de 20 chacune

48 emplacements (8 + 20 + 20)

Type de pipeline

Configuration

Emplacements requis

Sans état à un seul étage (source Kafka + sink)

maxPartitions = 8

8 emplacements

À deux étapes et avec état (source Kafka + brassage)

maxPartitions = 8, partitions de shuffle = 20

28 slots (8 + 20)

À trois étapes (source Kafka + shuffle + repartition)

maxPartitions = 8, deux étapes de brassage de 20 chacune

48 emplacements (8 + 20 + 20)

Si vous ne définissez pas maxPartitions, utilisez le nombre de partitions dans le sujet Kafka.

Ressources supplémentaires