Aller au contenu principal

Traitement asynchrone avec transformWithState (Bêta)

info

Bêta

Le traitement asynchrone pour l'API Python transformWithState basée sur les lignes est en version bêta. Consultez Versions préliminaires de Databricks.

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

Python transformWithState prend en charge le traitement asynchrone basé sur asyncio. En exécutant simultanément des opérations d'état et la logique utilisateur sur des clés de regroupement et en traitant par batch la communication inter-processus, le traitement asynchrone offre un throughput plus élevé que le traitement synchrone, avec seulement des modifications mineures du code. Ce gain de throughput ne nécessite aucune bibliothèque asynchrone tierce. Les utilisateurs avancés peuvent optimiser davantage leurs applications avec des modèles de programmation asynchrone et des bibliothèques compatibles avec l'asynchrone.

Pour utiliser le traitement asynchrone, implémentez un AsyncStatefulProcessor au lieu du StatefulProcessor synchrone. L’API AsyncStatefulProcessor reflète l’API synchrone StatefulProcessor, de sorte que la plupart des applications ne nécessitent que peu de modifications pour utiliser l’API asynchrone. Voir Mettre en œuvre un AsyncStatefulProcessor.

Pour l’API transformWithState synchrone et les concepts fondamentaux, consultez Créer une application avec état personnalisée avec transformWithState.

remarque

Le traitement asynchrone est disponible uniquement pour l’API transformWithState basée sur les lignes Python. Il n’est pas pris en charge pour transformWithStateInPandas ou pour l’API Scala transformWithState. Le traitement asynchrone n’est pas pris en charge dans le compute serverless.

Mettre en œuvre un AsyncStatefulProcessor

Pour convertir un StatefulProcessor synchrone en AsyncStatefulProcessor, apportez les modifications suivantes :

  • Définissez les méthodes d’API (init, close, handleInputRows, handleExpiredTimer et handleInitialState) avec le mot-clé async def.
  • Lisez et mettez à jour les valeurs d'état et de minuteur avec await, ou exécutez-les à l'aide de la bibliothèque asyncio de Python. Ceci s'applique aux opérations d'état telles que valueState.get() et aux opérations de minuteur 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 des données dans des variables membres ou dans des systèmes externes, Databricks vous recommande de réécrire la logique pour qu’elle soit sécurisée pour une exécution simultanée. Comme handleInputRows et handleExpiredTimer peuvent s’exécuter simultanément sur des clés de regroupement, les exécutions entrelacées ne doivent pas corrompre les données partagées. La plupart des applications répondent déjà à cette exigence.
  • Databricks recommande de ne pas intercepter ni supprimer les erreurs provenant des opérations d’état. Apache Spark gère ces erreurs pour vous. Si une opération d'état échoue, Apache Spark fait échouer la tâche et la relance.
    • Dans un AsyncStatefulProcessor, les erreurs d'opération d'état sont gérées pour vous et ne sont jamais transmises à votre code.
    • Dans un StatefulProcessor synchrone, les erreurs d'opération d'état sont générées dans votre code, mais leur suppression peut compromettre l'exactitude des données.

Exemple : compter les lignes pour chaque clé de regroupement

L’exemple suivant définit un AsyncCountProcessor qui compte le nombre de lignes pour chaque clé de regroupement. La variable value_schema définit le schéma du ValueState qui stocke le décompte en cours. Par rapport à un StatefulProcessor synchrone, les changements sont le mot-clé async def sur chaque méthode et await sur les opérations de lecture et de mise à jour de l’état. L’appel à getValueState dans init reste synchrone. Définissez le processeur comme dans le code suivant :

Python
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

Exécuter une query avec un processeur asynchrone

Pour exécuter une requête avec un processeur asynchrone, transmettez votre AsyncStatefulProcessor à transformWithState. La query utilise la même syntaxe que le chemin synchrone. Les APIs asynchrones et synchrones partagent le même format d'état ; vous pouvez donc basculer une query existante entre une AsyncStatefulProcessor et une StatefulProcessor synchrone tout en réutilisant le même point de contrôle.

Exemple : compter les événements dans le dataset exemple events

L'exemple suivant exécute AsyncCountProcessor sur le dataset d'exemple events. Chaque enregistrement possède un champ time (secondes epoch) et un champ action avec la valeur Open ou Close. La query effectue un regroupement par action et compte les événements pour chaque type d'action. Pour plus de datasets d'exemple, consultez Sample datasets.

La variable input_schema définit le schéma des enregistrements source, et la variable output_schema définit le schéma des lignes émises par le processeur. Pour lire le dataset exemple en tant que Stream, définissez les deux schémas, puis start la query comme dans le code suivant :

Python
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 query terminée, affichez le nombre en cours pour chaque type d'action comme dans le code suivant :

Python
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 minuteur

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

Le tableau suivant répertorie les opérations qui renvoient un résultat unique 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

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 répertorie les opérations qui renvoient 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

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érez avec async for comme dans le code suivant :

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

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

Pour une description de chaque type d’état, consultez Custom state types.

Optimiser avec des modèles de programmation asynchrone

Le traitement asynchrone est utile lorsque votre logique attend des opérations externes, telles que des 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 simultanées avec asyncio.gather

L’exemple suivant utilise asyncio.gather pour déclencher toutes les requêtes HTTP par ligne simultanément et attendre qu’elles se terminent, puis stocke le score maximal dans l’état. Définissez le processeur comme dans le code suivant :

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