Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
Langues prises en charge
Le mode en temps réel prend en charge Scala, Java et Python.
Types de calcul
Le mode en temps réel prend en charge les types de calcul suivants :
| Type de capacité de calcul | Soutenu |
|---|---|
| Dédié (anciennement : utilisateur unique) | ✓ |
| Standard (anciennement : partagé) | ✓ (uniquement Python) |
| Pipelines Lakeflow sur classique | Non pris en charge en tant que streaming structuré. Prise en charge par le biais de la configuration du pipeline. Consultez Utiliser le mode en temps réel dans les pipelines Lakeflow. |
| Pipelines Lakeflow sur serverless | Non pris en charge en tant que streaming structuré. Prise en charge par le biais de la configuration du pipeline. Consultez Utiliser le mode en temps réel dans les pipelines Lakeflow. |
| Sans serveur | Non pris en charge |
Pour les charges de travail sensibles à la latence avec des fonctions définies par l’utilisateur, Databricks recommande d’utiliser le mode d’accès dédié. Consultez les fonctions Table.
Modes d'exécution
Le mode en temps réel prend uniquement en charge le mode de mise à jour :
| Mode d’exécution | Soutenu |
|---|---|
| Mode de mise à jour | ✓ |
| Mode d’ajout | Non pris en charge |
| Mode complet | Non pris en charge |
Sources et puits
Le mode en temps réel prend en charge les sources et récepteurs suivants :
| Source ou récepteur | En tant que source | En tant que récepteur |
|---|---|---|
| Apache Kafka | ✓ | ✓ |
| Event Hubs (à l’aide du 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 |
Récepteurs arbitraires (utilisation de forEachWriter) |
Sans objet | ✓ |
Opérateurs
Le mode en temps réel prend en charge la plupart des opérateurs de streaming structuré :
Opérations sans état
| Opérateur | Soutenu |
|---|---|
| Sélection | ✓ |
| Projection | ✓ |
mapPartitions |
Non pris en charge (voir limitation) |
| Union | ✓ (avec certaines limitations) |
UDFs
| Opérateur | Soutenu |
|---|---|
| Scala UDF | ✓ (avec certaines limitations) |
| Python UDF | ✓ (avec certaines limitations) |
Aggregation
| Function | Soutenu |
|---|---|
| sum | ✓ |
| count | ✓ |
| max | ✓ |
| min | ✓ |
| avg | ✓ |
| Fonctions d’agrégation | ✓ |
Windowing
| Opérateur | Soutenu |
|---|---|
| Tumbling | ✓ |
| Sliding | ✓ |
| Session | Non pris en charge |
Deduplication
| Opérateur | Soutenu |
|---|---|
| dropDuplicates | ✓ |
| dropDuplicatesWithinWatermark | ✓ |
Diffuser en continu vers la jointure de table
| Opérateur | Soutenu |
|---|---|
| Jointure interne | ✓ |
| Jointure externe | ✓ |
| Jointure de table de diffusion (taille de table de 10 Mo ou moins) | ✓ |
| Jointure de table (sans diffusion) | Non pris en charge |
Diffuser en continu la jointure
| Opérateur | Soutenu |
|---|---|
| Jointure interne | ✓ (Databricks Runtime 18 LTS et versions ultérieures, avec certaines configurations) |
| Jointure externe | Non pris en charge |
Note
Pour utiliser le flux pour diffuser des jointures en mode temps réel, vous devez définir des configurations Spark supplémentaires. Pour plus d’informations sur les configurations et la configuration requise pour exécuter plusieurs flux, consultez Stream to stream joins.
Opérateur avec état arbitraire
| Opérateur | Soutenu |
|---|---|
| (plat)MapGroupsWithState | Non pris en charge |
| transformWithState | ✓ (avec certaines différences) |
Récepteurs définis par l’utilisateur
| Récepteur | Soutenu |
|---|---|
| forEach | ✓ |
| forEachBatch | Non pris en charge |
Considérations spéciales
Certains opérateurs et fonctionnalités ont 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 personnalisées avec état, Databricks prend en charge transformWithState, une API dans Apache Spark Structured Streaming. Consultez Créer une application avec état personnalisé pour plus d’informations sur l’API et les extraits de code.
Toutefois, l’API se comporte différemment en mode temps réel que dans les requêtes de micro-lots.
- Le mode en temps réel appelle la
handleInputRows(key: String, inputRows: Iterator[T], timerValues: TimerValues)méthode pour chaque ligne.- L’itérateur
inputRowsretourne une valeur unique. Le mode micro-lot l'appelle une fois pour chaque clé, et l'itérateurinputRowsrenvoie toutes les valeurs d'une clé dans le micro-lot. - Tenir compte de cette différence lors de l’écriture de votre code
- L’itérateur
- Les temporisateurs d’événements 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’API basée surtransformWithStateles lignes, qui utiliseRowdes objets plutôt que des dataFrames pandas. - En mode temps réel, les minuteurs sont retardés en fonction de l’arrivée des données :
- Si un minuteur est planifié 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 traitement de longue durée se termine, le minuteur se déclenche avant l’arrêt du traitement.
Note
Dans Databricks Runtime 18.1 et ci-dessous, si vous utilisez transformWithState et le mode en temps réel pour Python avec un débit faible, moins de 5 enregistrements par seconde, vous pouvez voir des latences accrues allant jusqu’à quelques centaines de millisecondes. Databricks recommande la mise à niveau vers Databricks Runtime 18.2 et versions ultérieures à résoudre.
Python fonctions définies par l’utilisateur 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 UDF | Soutenu |
|---|---|
| Python fonctions UDF scalaires (fonctions définies par l’utilisateur scalaires Python)) | ✓ |
| Fonctions UDF scalaires de flèche | ✓ |
| Fonctions UDF scalaires Pandas (fonctions définies par l’utilisateur pandas) | ✓ |
Fonction fléchée (mapInArrow) |
✓ |
| Pandas, fonction (Map) | ✓ |
Regroupement avec état (UDAF)
| Type UDF | Soutenu |
|---|---|
transformWithState (interface uniquement Row ) |
✓ |
transformWithStateInPandas |
Non pris en charge. Utilisez plutôt l’API basée sur transformWithState les lignes, qui utilise Row des objets plutôt que des dataFrames pandas. Voir transformWithStateInPandas non pris en charge pour plus d’informations . |
applyInPandasWithState |
Non pris en charge |
Regroupement non avec état (UDAF)
| Type UDF | Soutenu |
|---|---|
apply |
Non pris en charge |
applyInArrow |
Non pris en charge |
applyInPandas |
Non pris en charge |
Fonctions de table
| Type UDF | Soutenu |
|---|---|
| UDTF (Python fonctions de table définies par l’utilisateur (UDF)) | Non pris en charge |
| UC UDF | Non pris en charge |
Il existe plusieurs points à prendre en compte lors de l’utilisation de Python fonctions définies par l’utilisateur en mode temps réel :
- Pour réduire la latence, définissez la taille du lot flèche (
spark.sql.execution.arrow.maxRecordsPerBatch) sur 1.- Compromis : cette configuration optimise la latence au détriment du débit. Pour la plupart des charges de travail, ce paramètre est recommandé.
- Augmentez la taille du lot uniquement si un débit plus élevé est nécessaire pour prendre en charge le volume d’entrée, en acceptant l’augmentation potentielle de la latence.
- Les fonctions pandas et UDF ne fonctionnent pas correctement avec une taille de lot Arrow de 1.
- Si vous utilisez des fonctions UDFs ou pandas, définissez la taille du lot Arrow sur une valeur supérieure (par exemple, 100 ou plus).
- Cela implique une latence plus élevée. Databricks recommande d’utiliser une fonction ou une fonction UDF de flèche si possible.
-
transformWithStateInPandasn’est pas pris en charge en mode temps réel. Utilisez plutôt l’API basée surtransformWithStateles lignes, qui utiliseRowdes objets plutôt que des dataFrames pandas. ConsulteztransformWithStateInPandasnon pris en charge et des exemples de mode Real-time pour obtenir un exemple de Python opérationnel à l’aide de l’API basée sur les lignes. - Pour les charges de travail sensibles à la latence avec des fonctions définies par l’utilisateur, Databricks recommande d’utiliser le mode d’accès dédié. En mode d’accès standard, la surcharge d’isolation de la sécurité peut ralentir les performances des fonctions utilisateur.