Aller au contenu principal

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

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

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

Non applicable

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

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

mapPartitions

Non pris en charge (voir la limitation)

Union

✓ (avec certaines limitations)

Opérateur

Pris en charge

Sélection

Projection

mapPartitions

Non pris en charge (voir la limitation)

Union

✓ (avec certaines limitations)

UDF

Opérateur

Pris en charge

UDF Scala

✓ (avec certaines limitations)

Python UDF

✓ (avec certaines limitations)

Opérateur

Pris en charge

UDF Scala

✓ (avec certaines limitations)

Python UDF

✓ (avec certaines limitations)

Agrégation

Fonction

Pris en charge

somme

Décompte

max

min

moy.

Fonctions d'agrégation

Fonction

Pris en charge

somme

Décompte

max

min

moy.

Fonctions d'agrégation

Fenêtrage

Opérateur

Pris en charge

Tumbling

Défilement

Session

Non pris en charge

Opérateur

Pris en charge

Tumbling

Défilement

Session

Non pris en charge

Déduplication

Opérateur

Pris en charge

dropDuplicates

dropDuplicatesWithinWatermark

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

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

Opérateur

Pris en charge

Jointure intérieure

✓ (Databricks Runtime 18 et versions supérieures, avec certaines configurations)

Jointure externe

Non pris en charge

remarque

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

✓ (avec quelques différences)

Opérateur

Pris en charge

(flat)MapGroupsWithState

Non pris en charge

transformWithState

✓ (avec quelques différences)

Puits définis par l'utilisateur

Puits

Pris en charge

forEach

ForEachBatch

Non pris en charge

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 inputRows renvoie une valeur unique. Le mode micro-batch l'appelle une fois pour chaque clé, et l'itérateur inputRows renvoie toutes les valeurs pour une clé dans le micro batch.
    • Tenez compte de cette différence lors de l'écriture de votre code.
  • Les minuteurs d'heure d'événement ne sont pas pris en charge en mode temps réel.

  • transformWithStateInPandas n'est pas pris en charge en mode temps réel. Utilisez plutôt l’API transformWithState basée sur les lignes, qui utilise des objets Row plutô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.
remarque

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

Fonction Pandas (Carte)

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

Fonction Pandas (Carte)

Regroupement avec état (UDAF)

Type d'UDF

Pris en charge

transformWithState (seulement Row interface)

transformWithStateInPandas

Non pris en charge. Utilisez plutôt l'API transformWithState basée sur les lignes, qui utilise des objets Row plutôt que des DataFrames pandas. Consultez transformWithStateInPandas non pris en charge pour plus de détails.

applyInPandasWithState

Non pris en charge

Type d'UDF

Pris en charge

transformWithState (seulement Row interface)

transformWithStateInPandas

Non pris en charge. Utilisez plutôt l'API transformWithState basée sur les lignes, qui utilise des objets Row plutôt que des DataFrames pandas. Consultez transformWithStateInPandas non pris en charge pour plus de détails.

applyInPandasWithState

Non pris en charge

Regroupement sans état (UDAF)

Type d'UDF

Pris en charge

apply

Non pris en charge

applyInArrow

Non pris en charge

applyInPandas

Non pris en charge

Type d'UDF

Pris en charge

apply

Non pris en charge

applyInArrow

Non pris en charge

applyInPandas

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

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.
  • transformWithStateInPandas n'est pas pris en charge en mode temps réel. Utilisez plutôt l'API transformWithState basée sur les lignes, qui utilise des objets Row plutôt que des DataFrames pandas. Consultez transformWithStateInPandas non 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.