Repartitionnement d'état à la demande pour les query streaming avec état
Aperçu
Cette fonctionnalité est en aperçu public.
Le repartitionnement d'état à la demande vous permet de redimensionner le nombre de partitions pour une query Structured Streaming avec état sans perdre l'état du point de contrôle.
Sans repartitionnement de l'état à la demande, vous définissez le nombre de partitions de mélange lors de la création du point de contrôle. Si vous modifiez spark.sql.shuffle.partitions, les query avec des points de contrôle existants ignorent la nouvelle valeur. L'application d'un nouveau nombre de partitions nécessite que vous redémarriez la query avec un nouveau point de contrôle.
Le repartitionnement d'état à la demande présente les avantages suivants :
- Optimiser les queries en redimensionnant le nombre de partitions sans reconstruire le point de contrôle.
- Montez en charge les queries à la hausse ou à la baisse pour correspondre aux changements de charge de travail.
Exigences
- Databricks Runtime 18 et versions ultérieures.
- La requête doit utiliser le fournisseur de stockage d'état RocksDB. Sur DBR 17.3 ou version supérieure, RocksDB est le fournisseur de stockage d'état default. Consultez Configurer le stockage d'état RocksDB sur Databricks.
Modifier le nombre de partitions
Utilisez la configuration Spark spark.sql.streaming.stateStore.partitions et redémarrez la query pour modifier le nombre de partitions d'état de brassage et de streaming :
- Python
- Scala
query.stop()
spark.conf.set("spark.sql.streaming.stateStore.partitions", "<numPartitions>")
query = df.writeStream.start()
query.stop()
spark.conf.set("spark.sql.streaming.stateStore.partitions", "<numPartitions>")
val query = df.writeStream.start()
Pour les queries avec état, spark.sql.streaming.stateStore.partitions prévaut sur spark.sql.shuffle.partitions. Une fois que la query redémarre et que le dernier micro-batch planifié est terminé, la query exécute une opération de repartitionnement pour redistribuer les données d'état dans le nouveau nombre de partitions. Une fois l'opération de repartitionnement terminée, la query reprend le traitement.
Surveiller l'état de repartitionnement
Une fois le prochain micro-lot terminé, StreamingQueryProgress événements incluent la durée de l'opération de répartition. Dans les métriques durationMs d'un événement, controlBatch.REPARTITION affiche la valeur de durée en millisecondes. Des tailles d'état plus importantes pourraient augmenter le temps de repartitionnement. Consultez le monitoring des queries Structured Streaming sur Databricks.
Exemple de Structured Streaming
L'exemple suivant réduit une query de 200, la valeur par default, à 100 partitions de brassage. Arrêter la query, définir le nouveau nombre de partitions et redémarrer :
- Python
- Scala
# Start the query with the default partition count (200)
query = (df
.withWatermark("event_time", "10 minutes")
.groupBy(
window("event_time", "5 minutes"),
"id")
.count()
.writeStream
.format("delta")
.option("checkpointLocation", "/checkpoint/path")
.outputMode("append")
.start()
)
# Stop the query and scale down to 100 partitions
query.stop()
spark.conf.set("spark.sql.streaming.stateStore.partitions", "100")
# Restart the query with the same options
query = (df
.withWatermark("event_time", "10 minutes")
.groupBy(
window("event_time", "5 minutes"),
"id")
.count()
.writeStream
.format("delta")
.option("checkpointLocation", "/checkpoint/path")
.outputMode("append")
.start()
)
// Start the query with the default partition count (200)
val query = df
.withWatermark("event_time", "10 minutes")
.groupBy(
window($"event_time", "5 minutes"),
$"id")
.count()
.writeStream
.format("delta")
.option("checkpointLocation", "/checkpoint/path")
.outputMode("append")
.start()
// Stop the query and scale down to 100 partitions
query.stop()
spark.conf.set("spark.sql.streaming.stateStore.partitions", "100")
// Restart the query with the same options
val query2 = df
.withWatermark("event_time", "10 minutes")
.groupBy(
window($"event_time", "5 minutes"),
$"id")
.count()
.writeStream
.format("delta")
.option("checkpointLocation", "/checkpoint/path")
.outputMode("append")
.start()
LakeFlow Pipelines example
Dans les LakeFlow Pipelines, définissez spark.sql.streaming.stateStore.partitions à l'aide du paramètre spark_conf sur le décorateur @dp.table ou @dp.append_flow.
Définir des partitions sur un flux :
from pyspark import pipelines as dp
from pyspark.sql import functions as F
source_path = "/databricks-datasets/iot-stream/data-device/"
dp.create_streaming_table("target_table")
@dp.append_flow(
target="target_table",
name="my_flow_1",
spark_conf={"spark.sql.streaming.stateStore.partitions": "100"}
)
def my_flow_1():
return (spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load(source_path)
.withColumn("timestamp", F.to_timestamp("timestamp"))
.withWatermark("timestamp", "10 minutes")
.groupBy(F.window("timestamp", "5 minutes"), "id")
.count())
Définissez les partitions au niveau de la table pour le flux default :
from pyspark import pipelines as dp
from pyspark.sql import functions as F
source_path = "/databricks-datasets/iot-stream/data-device/"
@dp.table(
name="table_1",
spark_conf={"spark.sql.streaming.stateStore.partitions": "100"}
)
def table_1():
return (spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load(source_path)
.withColumn("timestamp", F.to_timestamp("timestamp"))
.withWatermark("timestamp", "10 minutes")
.groupBy(F.window("timestamp", "5 minutes"), "id")
.count())