Repartición de estado a petición para consultas de streaming con estado

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

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