Nota
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
Importante
Esta característica está en versión preliminar pública.
La repartición de estado a petición permite cambiar el tamaño del número de particiones de una consulta de streaming estructurado con estado sin perder el estado de punto de control.
Sin la repartición de estado a petición se establece el número de particiones aleatorias durante la creación del punto de control. Si cambia spark.sql.shuffle.partitions, las consultas con puntos de control existentes omiten el nuevo valor. La aplicación de un nuevo recuento de particiones requiere reiniciar la consulta con un nuevo punto de control.
La repartición de estados a petición tiene las siguientes ventajas:
- Ajuste las consultas cambiando el número de particiones sin volver a generar el punto de control.
- Escale o reduzca verticalmente las consultas para que coincidan con los cambios de carga de trabajo.
Requirements
- Databricks Runtime 18 LTS y versiones posteriores.
- La consulta debe usar el proveedor de almacén de estado de RocksDB. En DBR 17.3 o superior, RocksDB es el proveedor de almacén de estado predeterminado. Consulte Configuración de almacenes de estados de RocksDB en Azure Databricks.
Cambiar el número de particiones
Use la configuración de Spark spark.sql.streaming.stateStore.partitions y reinicie la consulta para cambiar el número de particiones de shuffle y del estado de streaming:
Python
query.stop()
spark.conf.set("spark.sql.streaming.stateStore.partitions", "<numPartitions>")
query = df.writeStream.start()
Scala
query.stop()
spark.conf.set("spark.sql.streaming.stateStore.partitions", "<numPartitions>")
val query = df.writeStream.start()
En el caso de las consultas con estado, spark.sql.streaming.stateStore.partitions tiene prioridad sobre spark.sql.shuffle.partitions. Una vez reiniciada la consulta y se completa la última microbatch planeada, la consulta ejecuta una operación de repartición para redistribuir los datos de estado en el nuevo número de particiones. Una vez completada la operación de repartición, la consulta reanuda su procesamiento.
Supervisión del estado de repartición
Una vez completado el siguiente microlote, los eventos StreamingQueryProgress incluyen la duración de la operación de repartición. En las métricas de durationMs un evento, controlBatch.REPARTITION muestra el valor de duración en milisegundos. Los tamaños de estado más grandes pueden aumentar el tiempo de repartición. Consulte Supervisión de consultas de streaming estructurado en Azure Databricks.
Ejemplo de Structured Streaming
El siguiente ejemplo reduce una consulta de 200, el valor predeterminado, a 100 particiones de mezcla. Detenga la consulta, establezca el nuevo recuento de particiones y reinicie:
Python
# 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()
)
Scala
// 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()
Ejemplo de canalizaciones de Lakeflow
En las canalizaciones de Lakeflow, establezca spark.sql.streaming.stateStore.partitions mediante el parámetro spark_conf en el decorador @dp.table o @dp.append_flow.
Establecer particiones en un flujo:
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())
Establezca particiones en el nivel de tabla para el flujo predeterminado:
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())