Aller au contenu principal

Optimiser le traitement avec état avec des filigranes

Pour gérer efficacement les données conservées dans l'état, utilisez des filigranes lors de l'exécution du traitement de Stream avec état dans les LakeFlow Pipelines, y compris les agrégations, les jointures et la déduplication. Les sections suivantes montrent comment appliquer des filigranes dans vos queries de pipeline, avec des exemples des opérations recommandées.

Pour en savoir plus sur les tables de streaming, consultez Tables de streaming.

remarque

Pour s'assurer que les queries qui effectuent des agrégations sont traitées de manière incrémentielle et non entièrement recalculées à chaque mise à jour, vous devez utiliser des watermarks.

Qu'est-ce qu'un filigrane ?

Dans le traitement de flux, un *filigrane* est une fonctionnalité Apache Spark qui peut définir un threshold temporel pour le traitement des données lors de l'exécution d'opérations avec état telles que les agrégations. Les données entrantes sont traitées jusqu'à ce que le threshold soit atteint, moment auquel la fenêtre de temps définie par le threshold est fermée. Les filigranes peuvent être utilisés pour éviter les problèmes pendant le traitement des requêtes, principalement lors du traitement de datasets plus volumineux ou de traitements de longue durée. Ces problèmes peuvent inclure une latence élevée dans la production des résultats et même des erreurs de mémoire insuffisante (OOM) en raison de la quantité de données conservées en état pendant le traitement. Parce que les données de streaming sont intrinsèquement désordonnées, les filigranes supportent également le calcul correct d'opérations comme les agrégations par fenêtre temporelle.

Pour en savoir plus sur l'utilisation des filigranes dans le traitement de flux, consultez Watermarking in Apache Spark Structured Streaming et Appliquer des filigranes pour contrôler les thresholds de traitement des données.

Comment définissez-vous un filigrane ?

Vous définissez un watermark en spécifiant un champ Timestamp et une valeur représentant le threshold de temps pour l'arrivée des données tardives . Les données sont considérées comme tardives si elles arrivent après le threshold de temps défini. Par exemple, si le threshold est défini sur 10 minutes, les enregistrements arrivant après le threshold de 10 minutes pourraient être ignorés.

Étant donné que les enregistrements qui arrivent après le threshold défini peuvent être supprimés, il est important de sélectionner un threshold qui répond à vos exigences en matière de latence et d'exactitude. Le choix d'un threshold plus petit permet d'émettre les enregistrements plus tôt, mais signifie également que les enregistrements tardifs sont plus susceptibles d'être supprimés. Un threshold plus élevé signifie une attente plus longue, mais éventuellement une plus grande exhaustivité des données. En raison de la taille d'état plus importante, un threshold plus élevé peut également nécessiter des Ressources de calcul supplémentaires. Étant donné que la valeur du threshold dépend de vos données et de vos exigences de traitement, il est important de tester et de monitoring votre traitement pour déterminer un threshold optimal.

Vous utilisez la fonction withWatermark() en Python pour définir un filigrane. En SQL, utilisez la clause WATERMARK pour définir un watermark :

Python
withWatermark("timestamp", "3 minutes")

Utilisez les filigranes avec les jointures Stream-Stream

Pour les jointures de Stream à Stream, vous devez définir un filigrane des deux côtés de la jointure et une clause d'intervalle de temps. Étant donné que chaque source de jointure a une vue incomplète des données, la clause d'intervalle de temps est nécessaire pour indiquer au moteur de streaming quand aucune correspondance supplémentaire ne peut être établie. La clause d'intervalle de temps doit utiliser les mêmes champs que ceux utilisés pour définir les filigranes.

Comme il peut arriver que chaque Stream nécessite des thresholds différents pour les filigranes, les Streams n'ont pas besoin d'avoir les mêmes thresholds. Pour éviter de perdre des données, le moteur de streaming maintient un filigrane global basé sur le stream le plus lent.

L'exemple suivant joint un stream d'impressions publicitaires et un stream de clics d'utilisateurs sur des publicités. Dans cet exemple, un clic doit se produire dans les 3 minutes suivant l'impression. Une fois l'intervalle de temps de 3 minutes écoulé, les lignes de l'état qui ne peuvent plus être mises en correspondance sont supprimées.

Python
from pyspark import pipelines as dp

dp.create_streaming_table("adImpressionClicks")
@dp.append_flow(target = "adImpressionClicks")
def joinClicksAndImpressions():
clicksDf = (read_stream("rawClicks")
.withWatermark("clickTimestamp", "3 minutes")
)
impressionsDf = (read_stream("rawAdImpressions")
.withWatermark("impressionTimestamp", "3 minutes")
)
joinDf = impressionsDf.alias("imp").join(
clicksDf.alias("click"),
expr("""
imp.userId = click.userId AND
clickAdId = impressionAdId AND
clickTimestamp >= impressionTimestamp AND
clickTimestamp <= impressionTimestamp + interval 3 minutes
"""),
"inner"
).select("imp.userId", "impressionAdId", "clickTimestamp", "impressionSeconds")

return joinDf

Effectuer des agrégations fenêtrées avec des filigranes

Une opération avec état courante sur les données de streaming est une agrégation fenêtrée. Les agrégations fenêtrées sont similaires aux agrégations groupées, sauf que les valeurs agrégées sont renvoyées pour l'ensemble des lignes qui font partie de la fenêtre définie.

Une fenêtre peut être définie comme une certaine longueur, et une opération d'agrégation peut être effectuée sur toutes les lignes qui font partie de cette fenêtre. Spark Streaming prend en charge trois types de fenêtres :

  • **Fenêtres basculantes (fixes)** : Une série d'intervalles de temps de taille fixe, non chevauchants et contigus. Un enregistrement d'entrée appartient à une seule fenêtre.
  • Fenêtres glissantes : Semblables aux fenêtres basculantes, les fenêtres glissantes sont de taille fixe, mais les fenêtres peuvent se chevaucher, et un enregistrement peut se retrouver dans plusieurs fenêtres.

Lorsque les données arrivent après la fin de la fenêtre plus la longueur du filigrane, aucune nouvelle donnée n'est acceptée pour la fenêtre, le résultat de l'agrégation est émis et l'état de la fenêtre est abandonné.

L'exemple suivant calcule une somme d'impressions toutes les 5 minutes à l'aide d'une fenêtre fixe. Dans cet exemple, la clause select utilise l'alias impressions_window, puis la fenêtre elle-même est définie comme faisant partie de la clause GROUP BY. La fenêtre doit être basée sur la même colonne Timestamp que le filigrane, la colonne clickTimestamp dans cet exemple.

SQL
CREATE OR REFRESH STREAMING TABLE
gold.adImpressionSeconds
AS SELECT
impressionAdId, window(clickTimestamp, "5 minutes") as impressions_window, sum(impressionSeconds) as totalImpressionSeconds
FROM STREAM
(silver.adImpressionClicks)
WATERMARK
clickTimestamp DELAY OF INTERVAL 3 MINUTES
GROUP BY
impressionAdId, window(clickTimestamp, "5 minutes")

Un exemple similaire en Python pour calculer le profit sur des fenêtres horaires fixes :

Python
from pyspark import pipelines as dp

@dp.table()
def profit_by_hour():
return (
spark.readStream.table("sales")
.withWatermark("timestamp", "1 hour")
.groupBy(window("timestamp", "1 hour").alias("time"))
.aggExpr("sum(profit) AS profit")
)

Dédupliquer les enregistrements de streaming

Structured Streaming offre des garanties de traitement « exactement une fois », mais ne dédouble pas automatiquement les enregistrements des sources de données. Par exemple, étant donné que de nombreuses files d'attente de messages offrent des garanties d'au moins une livraison, des enregistrements en double doivent être attendus lors de la lecture à partir de l'une de ces files d'attente de messages. Vous pouvez utiliser la fonction dropDuplicatesWithinWatermark() pour dédupliquer les enregistrements sur n'importe quel champ spécifié, en supprimant les doublons d'un Stream, même si certains champs diffèrent (tels que l'heure de l'événement ou l'heure d'arrivée). Vous devez spécifier un filigrane pour utiliser la fonction dropDuplicatesWithinWatermark(). Toutes les données en double qui arrivent dans la plage de temps spécifiée par le filigrane sont supprimées.

Les données ordonnées sont importantes car les données désordonnées font avancer la valeur du watermark de manière incorrecte. Ensuite, lorsque des données plus anciennes arrivent, elles sont considérées comme tardives et supprimées. Utilisez l'option withEventTimeOrder pour traiter le snapshot initial dans l'ordre, en fonction du Timestamp spécifié dans le watermark. L'option withEventTimeOrder peut être déclarée dans le code définissant le dataset ou dans les paramètres du pipeline à l'aide de spark.databricks.delta.withEventTimeOrder.enabled. Par exemple :

JSON
{
"spark_conf": {
"spark.databricks.delta.withEventTimeOrder.enabled": "true"
}
}
remarque

L'option withEventTimeOrder est prise en charge uniquement avec Python.

Dans l’exemple suivant, les données sont traitées dans l’ordre clickTimestamp, et les enregistrements arrivant à 5 secondes d’intervalle qui contiennent des colonnes userId et clickAdId en double sont supprimés.

Python
clicksDedupDf = (
spark.readStream.table
.option("withEventTimeOrder", "true")
.table("rawClicks")
.withWatermark("clickTimestamp", "5 seconds")
.dropDuplicatesWithinWatermark(["userId", "clickAdId"]))

Optimiser la configuration du pipeline pour le traitement avec état

Afin de prévenir les problèmes de production et la latence excessive, Databricks recommande d'activer la gestion de l'état basée sur RocksDB pour votre traitement de Stream avec état, particulièrement si votre traitement nécessite l'enregistrement d'une grande quantité d'état intermédiaire.

Les pipelines Serverless gèrent automatiquement les configurations de stockage d'état.

Vous pouvez activer la gestion d'état basée sur RocksDB en définissant la configuration suivante avant de déployer un pipeline :

JSON
{
"configuration": {
"spark.sql.streaming.stateStore.providerClass": "com.databricks.sql.streaming.state.RocksDBStateStoreProvider"
}
}

Pour en savoir plus sur le magasin d'état RocksDB, y compris les recommandations de configuration pour RocksDB, consultez Configurer le magasin d'état RocksDB sur Databricks.

Pour les opérations avec état qui nécessitent une latence de l'ordre de la milliseconde, consultez Utiliser le mode temps réel dans les LakeFlow Pipelines.