Référence de mode en temps réel
Langues prises en charge
Le mode temps réel prend en charge Scala, Java et Python.
Types de compute
Le mode temps réel prend en charge les types de compute suivants :
Type de compute | Pris en charge |
|---|---|
Dédié (anciennement : utilisateur unique) | ✓ |
Standard (anciennement : partagé) | ✓ (Python uniquement) |
LakeFlow Pipelines on classic | Non pris en charge comme Structured Streaming. Pris en charge par la configuration du pipeline. Consultez Utiliser le mode temps réel dans LakeFlow Pipelines. |
LakeFlow Pipelines sur Serverless | Non pris en charge comme Structured Streaming. Pris en charge par la configuration du pipeline. Consultez Utiliser le mode temps réel dans LakeFlow Pipelines. |
Serverless | Non pris en charge |
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.
Modes d’exécution
Le mode temps réel prend en charge le mode de mise à jour uniquement :
Mode d'exécution | Pris en charge |
|---|---|
Mode de mise à jour | ✓ |
Mode d'ajout | Non pris en charge |
Mode complet | Non pris en charge |
Sources et puits
Le mode temps réel prend en charge les sources et dépôts suivants :
Source ou récepteur | Comme source | En tant que puits |
|---|---|---|
Apache Kafka | ✓ | ✓ |
Event Hubs (utilisant le connecteur Kafka) | ✓ | ✓ |
Kinesis | ✓ (mode EFO uniquement) | Non pris en charge |
AWS MSK | ✓ | Non pris en charge |
Delta | Non pris en charge | Non pris en charge |
Google Pub/Sub | Non pris en charge | Non pris en charge |
Apache Pulsar | Non pris en charge | Non pris en charge |
Puits arbitraires (utilisant | Non applicable | ✓ |
Opérateurs
Le mode en temps réel prend en charge la plupart des opérateurs Structured Streaming :
Opérations sans état
Opérateur | Pris en charge |
|---|---|
Sélection | ✓ |
Projection | ✓ |
| Non pris en charge (voir la limitation) |
Union |
UDF
Opérateur | Pris en charge |
|---|---|
UDF Scala | |
Python UDF |
Agrégation
Fonction | Pris en charge |
|---|---|
somme | ✓ |
Décompte | ✓ |
max | ✓ |
min | ✓ |
moy. | ✓ |
✓ |
Fenêtrage
Opérateur | Pris en charge |
|---|---|
Tumbling | ✓ |
Défilement | ✓ |
Session | Non pris en charge |
Déduplication
Opérateur | Pris en charge |
|---|---|
dropDuplicates | ✓ |
dropDuplicatesWithinWatermark | ✓ |
Stream to table join
Opérateur | Pris en charge |
|---|---|
Jointure intérieure | ✓ |
Jointure externe | ✓ |
Jointure de table broadcast (taille de table de 10 Mo ou moins) | ✓ |
Jointure de table (sans diffusion) | Non pris en charge |
Jointure stream à stream
Opérateur | Pris en charge |
|---|---|
Jointure intérieure | ✓ (Databricks Runtime 18 et versions supérieures, avec certaines configurations) |
Jointure externe | Non pris en charge |
Pour utiliser les jointures stream à stream en mode temps réel, vous devez définir des configurations Spark supplémentaires. Pour plus d'informations sur les configurations et les exigences pour exécuter plusieurs Streams, consultez Jointures de Stream à Stream.
Opérateur avec état arbitraire
Opérateur | Pris en charge |
|---|---|
(flat)MapGroupsWithState | Non pris en charge |
transformWithState |
Puits définis par l'utilisateur
Puits | Pris en charge |
|---|---|
forEach | ✓ |
ForEachBatch | Non pris en charge |
Considérations spéciales
Certains opérateurs et fonctionnalités présentent des considérations ou des différences spécifiques lorsqu'ils sont utilisés en mode temps réel.
transformWithState en mode temps réel
Pour la création d'applications avec état personnalisées, Databricks prend en charge transformWithState, une API dans Apache Spark Structured Streaming. Consultez Créer une application personnalisée avec état pour plus d’informations sur l’API et les extraits de code.
Cependant, l'API se comporte différemment en mode temps réel que dans les requêtes par micro-batch.
-
Le mode en temps réel appelle la méthode
handleInputRows(key: String, inputRows: Iterator[T], timerValues: TimerValues)pour chaque ligne.- L'itérateur
inputRowsrenvoie une valeur unique. Le mode micro-batch l'appelle une fois pour chaque clé, et l'itérateurinputRowsrenvoie toutes les valeurs pour une clé dans le micro batch. - Tenez compte de cette différence lors de l'écriture de votre code.
- L'itérateur
-
Les minuteurs d'heure d'événement ne sont pas pris en charge en mode temps réel.
-
transformWithStateInPandasn'est pas pris en charge en mode temps réel. Utilisez plutôt l’APItransformWithStatebasée sur les lignes, qui utilise des objetsRowplutôt que des DataFrames pandas. -
En mode temps réel, les temporisateurs se déclenchent avec un délai qui dépend de l'arrivée des données :
- Si un minuteur est programmé pour 10:00:00 mais qu'aucune donnée n'arrive, le minuteur ne se déclenche pas immédiatement.
- Si les données arrivent à 10:00:10, le minuteur se déclenche avec un délai de 10 secondes.
- Si aucune donnée n'arrive et que le batch de longue durée se termine, le minuteur se déclenche avant que le batch ne se termine.
Dans Databricks Runtime 18.1 et versions antérieures, si vous utilisez transformWithState et le mode temps réel pour Python avec un faible throughput, moins de 5 enregistrements par seconde, vous pourriez constater une augmentation des latences allant jusqu’à quelques centaines de millisecondes. Databricks recommande de passer à Databricks Runtime 18.2 et versions ultérieures pour résoudre le problème.
UDF Python en mode temps réel
Databricks prend en charge la majorité des fonctions définies par l'utilisateur (UDF) Python en mode temps réel :
Sans état
Type d'UDF | Pris en charge |
|---|---|
UDF scalaire Python (fonctions scalaires définies par l'utilisateur (UDF) Python) | ✓ |
UDF scalaire Arrow. | ✓ |
UDF scalaire Pandas (fonctions définies par l'utilisateur Pandas) | ✓ |
Fonction fléchée ( | ✓ |
Fonction Pandas (Carte) | ✓ |
Regroupement avec état (UDAF)
Type d'UDF | Pris en charge |
|---|---|
| ✓ |
| Non pris en charge. Utilisez plutôt l'API |
| Non pris en charge |
Regroupement sans état (UDAF)
Type d'UDF | Pris en charge |
|---|---|
| Non pris en charge |
| Non pris en charge |
| Non pris en charge |
Fonctions de table
Type d'UDF | Pris en charge |
|---|---|
UDTF (fonctions de table définies par l'utilisateur (UDTF) Python) | Non pris en charge |
UDF UC | Non pris en charge |
Plusieurs points sont à considérer lors de l'utilisation des UDF Python en mode temps réel :
-
Pour minimiser la latence, définissez la taille du batch Arrow (
spark.sql.execution.arrow.maxRecordsPerBatch) sur 1.- Compromis : cette configuration optimise la latence au détriment du throughput. Pour la plupart des workloads, ce paramètre est recommandé.
- Augmentez la taille du batch uniquement si un throughput plus élevé est nécessaire pour s'adapter au volume d'entrée, en acceptant l'augmentation potentielle de la latence.
-
Les UDF et fonctions Pandas donnent de mauvais résultats avec une taille de batch Arrow de 1.
- Si vous utilisez des UDF pandas ou des fonctions, définissez la taille de batch Arrow sur une valeur plus élevée (par exemple, 100 ou plus).
- Cela implique une latence plus élevée. Databricks recommande d'utiliser une UDF Arrow ou une fonction si possible.
-
transformWithStateInPandasn'est pas pris en charge en mode temps réel. Utilisez plutôt l'APItransformWithStatebasée sur les lignes, qui utilise des objetsRowplutôt que des DataFrames pandas. ConsulteztransformWithStateInPandasnon pris en charge et Exemples de mode temps réel pour un exemple Python fonctionnel utilisant l'API basée sur les lignes. -
Pour les charges de travail sensibles à la latence avec des UDF, Databricks vous recommande d'utiliser le mode d'accès dédié. En mode d’accès standard, le surcoût de l'isolation de sécurité peut ralentir les performances des UDF.