Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Important
Эта функция доступна в общедоступной предварительной версии.
Переразбиение состояния по запросу позволяет изменять количество разделов для запроса Structured Streaming с сохранением состояния без потери состояния контрольных точек.
Без переразбиения состояния по запросу количество разделов shuffle задаётся при создании контрольной точки. При изменении spark.sql.shuffle.partitionsзапросы с существующими контрольными точками игнорируют новое значение. Чтобы применить новое количество разделов, необходимо перезапустить запрос, используя новую контрольную точку.
Перераспределение состояния по запросу имеет следующие преимущества:
- Настройте запросы, изменив количество секций без перестроения контрольной точки.
- Масштабируйте запросы вверх или вниз, чтобы соответствовать изменениям рабочей нагрузки.
Requirements
- Databricks Runtime 18 LTS и выше.
- Запрос должен использовать RocksDB в качестве поставщика хранилища состояний. В DBR 17.3 или более поздней версии RocksDB является поставщиком хранилища состояний по умолчанию. См. статью Настройка хранилища состояний RocksDB в Azure Databricks.
Изменение количества секций
Используйте конфигурационный параметр Spark spark.sql.streaming.stateStore.partitions и перезапустите запрос, чтобы изменить количество разделов перемешивания и разделов состояния потоковой передачи:
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()
Для запросов с сохранением состояния spark.sql.streaming.stateStore.partitions имеет приоритет над spark.sql.shuffle.partitions. После перезапуска запроса и завершения последнего запланированного микропакета запрос выполняет операцию репартиционирования для перераспределения данных состояния по новому числу разделов. После завершения операции переразбиения обработка запроса возобновляется.
Отслеживание состояния переразбиения
После завершения следующего микропакета события StreamingQueryProgress содержат сведения о длительности операции переразбиения. В метриках durationMs события controlBatch.REPARTITION отображается значение длительности в миллисекундах. Более крупные размеры состояний могут увеличить время повторного фрагментирования. См. Мониторинг запросов структурированного потокового вещания на Azure Databricks.
Пример структурированной потоковой передачи
В следующем примере количество разделов shuffle для запроса уменьшается со значения по умолчанию 200 до 100. Остановите запрос, задайте новое число секций и перезапустите:
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()
Пример конвейеров Lakeflow
В конвейерах Lakeflow задайте spark.sql.streaming.stateStore.partitions с помощью параметра spark_conf в декораторе @dp.table или @dp.append_flow.
Установите разделы в потоке:
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())
Задайте секции на уровне таблицы для потока по умолчанию:
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())