Traitement asynchrone avec transformWithState (Beta)

Important

Le traitement asynchrone pour l’API Python basée transformWithState sur les lignes est en version Bêta. Consultez Azure Databricks versions préliminaires.

Le traitement asynchrone est disponible dans Databricks Runtime 19 et versions supérieures.

Python transformWithState prend en charge le traitement asynchrone construit sur asyncio. En exécutant simultanément les opérations d’état et la logique utilisateur entre les clés de groupement et en regroupant la communication inter-processus, le traitement asynchrone a un débit plus élevé que le traitement synchrone, avec seulement des modifications mineures de code. Ce gain de débit ne nécessite aucune bibliothèque asynchrone tierce. Les utilisateurs avancés peuvent optimiser davantage leurs applications grâce à des schémas de programmation asynchrones et des bibliothèques compatibles async.

Pour utiliser un traitement asynchrone, implémentez un AsyncStatefulProcessor au lieu du synchrone StatefulProcessor. L’API AsyncStatefulProcessor reflète l’API synchrone StatefulProcessor , donc la plupart des applications nécessitent seulement de petits changements pour utiliser l’API asynchrone. Voir Implémenter un AsyncStatefulProcessor.

Pour l’API synchrone transformWithState et les concepts de base, voir Construire une application avec état personnalisée avec transformWithState.

Note

Le traitement asynchrone est disponible uniquement pour l’API Python basée transformWithState sur des lignes. Il n’est pas pris en charge ni transformWithStateInPandas pour l’API Scala transformWithState . Le traitement asynchrone n’est pas pris en charge dans le calcul sans serveur.

Implémentez un AsyncStatefulProcessor

Pour convertir un synchrone StatefulProcessor en AsyncStatefulProcessorun , effectuez les modifications suivantes :

  • Définissez les méthodes API (init, close, handleInputRows, handleExpiredTimer, et handleInitialState) avec le async def mot-clé.
  • Lisez et mettez à jour les valeurs d'état et de minuterie avec await, ou exécutez-les en utilisant la asyncio bibliothèque Python. Cela s’applique aux opérations d’état telles que valueState.get() et aux opérations de minuterie telles que registerTimer. La création d’objets d’état, tels que handle.getValueState, reste synchrone.

Les considérations suivantes s’appliquent au traitement asynchrone :

  • Si votre application stocke les données dans des variables membres ou dans des systèmes externes, Databricks recommande de réécrire la logique pour qu’elle soit sûre pour une exécution concurrente. Parce que handleInputRows et handleExpiredTimer peuvent s’exécuter simultanément à travers les clés de groupement, les exécutions entrelacées ne doivent pas corrompre les données partagées. La plupart des candidatures remplissent déjà cette exigence.
  • Databricks recommande de ne pas détecter ni supprimer les erreurs issues des opérations d’état. Apache Spark gère ces erreurs pour vous. Si une opération d’état échoue, Apache Spark échoue à la tâche et la réessaie.
    • Dans un AsyncStatefulProcessor, les erreurs d’opérations d’état sont gérées pour vous et ne sont jamais intégrées à votre code.
    • Dans un système synchrone StatefulProcessor, des erreurs d’opération d’état sont générées dans votre code, mais les supprimer peut compromettre la correction des données.

Exemple : compter les lignes pour chaque clé de groupe

L’exemple suivant définit un AsyncCountProcessor qui compte le nombre de lignes pour chaque clé de groupement. La value_schema variable définit le schéma de le ValueState qui stocke le décompte courant. Comparé à un système synchrone StatefulProcessor, les changements sont le async def mot-clé de chaque méthode et await les opérations de lecture et de mise à jour d’état. L’appel à getValueState in init reste synchrone. Définissons le processeur comme dans le code suivant :

from pyspark.sql import Row
from pyspark.sql.streaming import AsyncStatefulProcessor, AsyncStatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, LongType

value_schema = StructType([StructField("count", LongType(), True)])

class AsyncCountProcessor(AsyncStatefulProcessor):
  async def init(self, handle: AsyncStatefulProcessorHandle) -> None:
    self.count = handle.getValueState("count", value_schema)

  async def handleInputRows(self, key, rows, timerValues):
    total = (await self.count.get() or (0,))[0]
    for _ in rows:
      total += 1
    await self.count.update((total,))
    yield Row(action=key[0], count=total)

  async def close(self) -> None:
    pass

Lancer une requête avec un processeur asynchrone

Pour lancer une requête avec un processeur asynchrone, passez votre AsyncStatefulProcessor à transformWithState. La requête utilise la même syntaxe que le chemin synchrone. Les API asynchrone et synchrone partagent le même format d’état, vous pouvez donc basculer une requête existante entre an AsyncStatefulProcessor et une synchrone StatefulProcessor tout en réutilisant le même point de contrôle.

Exemple : compter les événements dans l’ensemble events de données échantillon

L’exemple suivant s’applique AsyncCountProcessor à l’ensemble events de données échantillon. Chaque enregistrement possède un time champ (secondes d’époque) et un action champ avec la valeur Open ou Close. La requête regroupe par action et compte les événements pour chaque type d’action. Pour plus d’exemples de jeux de données, voir Exemples de jeux de données.

La input_schema variable définit le schéma des enregistrements sources, et la output_schema variable définit le schéma des lignes émises par le processeur. Pour lire l’échantillon de jeu de données comme un flux, définissez les deux schémas, puis lancez la requête comme dans le code suivant :

from pyspark.sql.types import StructType, StructField, StringType, LongType

input_schema = StructType([
  StructField("time", LongType(), True),
  StructField("action", StringType(), True),
])

output_schema = StructType([
  StructField("action", StringType(), True),
  StructField("count", LongType(), True),
])

events = (
  spark.readStream.schema(input_schema)
    .option("maxFilesPerTrigger", 10)
    .json("/databricks-datasets/structured-streaming/events")
)

q = (
  events.groupBy("action")
    .transformWithState(
      statefulProcessor=AsyncCountProcessor(),
      outputStructType=output_schema,
      outputMode="Update",
      timeMode="None",
    )
    .writeStream.format("memory")
    .queryName("async_counts")
    .trigger(availableNow=True)
    .start()
)

q.awaitTermination()

Une fois la requête terminée, consultez le nombre de joueurs pour chaque type d’action comme indiqué dans le code suivant :

display(spark.sql("SELECT action, MAX(count) AS count FROM async_counts GROUP BY action ORDER BY action"))

Opérations d’état asynchrone et de minuterie

Dans un AsyncStatefulProcessor, les opérations de variable d’état et de minuterie qui lisent ou écrivent des valeurs sont asynchrones. La plupart de ces opérations retournent un seul résultat que vous récupérez avec await. Les opérations qui retournent une collection renvoient plutôt un itérateur asynchrone que vous consommez avec async for. Pour une introduction aux async/await itérateurs asynchrones en Python, voir la documentation Python asyncio.

Le tableau suivant liste les opérations qui retournent un seul résultat que vous pouvez récupérer avec await:

Classe Opérations utilisant await
AsyncValueState exists, get, update, clear
AsyncMapState exists, getValue, containsKey, updateValue, removeKey, clear
AsyncListState exists, put, appendValue, appendList, clear
AsyncStatefulProcessorHandle registerTimer, deleteTimer

Le tableau suivant liste les opérations qui retournent un itérateur asynchrone que vous pouvez récupérer avec async for:

Classe Opérations utilisant async for
AsyncMapState iterator, keys, values
AsyncListState get
AsyncStatefulProcessorHandle listTimers

Exemple : async for

Par exemple, pour lire les valeurs dans un AsyncListState, itère avec async for comme dans le code suivant :

total = 0
async for value in self.items.get():
  total += value[0]

Les méthodes qui créent les objets d’état et suppriment les variables d’état restent synchrones : getValueState, getMapState, getListState, et deleteIfExists.

Pour une description de chaque type d’état, voir Types d’état personnalisés.

Optimiser avec des schémas de programmation asynchrones

Le traitement asynchrone est utile lorsque votre logique attend des opérations externes, telles que les requêtes réseau. Au lieu d’attendre chaque requête en séquence, utilisez asyncio pour exécuter les requêtes simultanément et réduire le temps d’inactivité.

Exemple : exécuter des requêtes concurrentes avec asyncio.gather

L’exemple suivant sert asyncio.gather à lancer simultanément toutes les requêtes HTTP par ligne et à attendre qu’elles se terminent, puis stocker le score maximal dans l’état. Définissons le processeur comme dans le code suivant :

import asyncio
import aiohttp
from pyspark.sql import Row
from pyspark.sql.streaming import AsyncStatefulProcessor

class HttpScoreRowGatherProcessor(AsyncStatefulProcessor):
  async def init(self, handle):
    self._score_state = handle.getValueState("last_score", "score double")
    self._session = aiohttp.ClientSession()

  async def _fetch_score(self, row) -> float:
    async with self._session.get(
      f"https://api.example.com/score/{row.event_id}"
    ) as resp:
      return (await resp.json())["score"]

  async def handleInputRows(self, key, rows, timerValues):
    user_id = key[0]
    scores = await asyncio.gather(*[self._fetch_score(row) for row in rows])

    max_score = max(scores)
    await self._score_state.update((max_score,))
    yield Row(user_id=user_id, score=max_score)

  async def close(self):
    await self._session.close()