Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 08.07 · 22 мин
Продвинутый
transformWithStateStatefulProcessorArbitrary State API v2ValueStateListStateMapStateTimerStateTTLRocksDB State Store

transformWithState: Arbitrary State API v2

В предыдущем уроке мы разобрали mapGroupsWithState и flatMapGroupsWithState — legacy-API для произвольной stateful-обработки. В Spark 4.0 (GA, 2025) появился их преемник: оператор transformWithState, который Apache официально называет Arbitrary State API v2. Это не косметический рефакторинг, а полностью переработанный API: объектно-ориентированная декларация состояний, композитные state-типы, декларативный TTL, несколько таймеров на ключ и нативный паритет PySpark/Scala.

В этом уроке мы разберём, чем v2 принципиально отличается от v1, как устроен StatefulProcessor, какие state-типы доступны и как мигрировать с mapGroupsWithState.

Зачем понадобился новый API?

Legacy mapGroupsWithState работает с 2017 года и в production покрывает базовые сценарии sessionization. Но за годы накопились ограничения, которые упирались в саму форму API:

  • Один GroupState[S] на ключ. Нельзя завести две независимые переменные состояния (например, lastSeen и eventBuffer) — приходится паковать всё в один case class и сериализовать целиком на каждое чтение/запись.
  • Один таймер на ключ. setTimeoutDuration или setTimeoutTimestamp — только одно значение. Если нужно «через 5 минут отправить heartbeat и через 30 минут закрыть сессию», приходится самому мультиплексировать таймеры в коде.
  • Нет TTL. Просрочку каждой записи нужно отслеживать вручную, проверяя state.hasTimedOut.
  • Нет композитных типов. Нет нативного ListState или MapState — список или мапу нужно хранить внутри объекта целиком, что убивает производительность на больших коллекциях.
  • Слабый Python-паритет. В PySpark вместо настоящего mapGroupsWithStateapplyInPandasWithState, который привязан к pandas-интерфейсу и медленнее Scala-версии.
  • Нет state schema evolution. Изменил поле в case class — получи IllegalStateException при рестарте из чекпоинта.
  • Логика смешана. Обработка новых строк и обработка timeout живут в одной функции, и приходится разветвлять её через if state.hasTimedOut: ....

Arbitrary State API v2 (SPARK-45654 umbrella, реализация в SPARK-47960, SPARK-50194, SPARK-51411, SPARK-51814 и десятках связанных тикетов) решает эти проблемы за счёт перехода от «функции обновления» к объекту-процессору с явным жизненным циклом.

INFO

Версии. Оператор transformWithState стал частью Apache Spark 4.0.0 (GA, июнь 2025). В 4.1 (preview/GA конец 2025) добавлены state schema evolution через Avro-кодирование, чейнинг stateful-операторов, расширенный state data source и улучшения батчевого режима. На Databricks transformWithState доступен с DBR 16.2+ и используется по умолчанию в Runtime 17.3+.

Концепции transformWithState

Вместо одной функции с тремя аргументами (key, events, state) v2-API предлагает определить класс, наследующий StatefulProcessor[K, I, O]:

import org.apache.spark.sql.streaming._

class MyProcessor extends StatefulProcessor[String, InputRow, OutputRow] {
  // 1. Объявление state-переменных через handle
  @transient private var counter: ValueState[Long] = _

  override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
    counter = getHandle.getValueState[Long]("counter", Encoders.scalaLong)
  }

  override def handleInputRows(
    key: String,
    inputRows: Iterator[InputRow],
    timerValues: TimerValues
  ): Iterator[OutputRow] = {
    val current = if (counter.exists()) counter.get() else 0L
    val updated = current + inputRows.size
    counter.update(updated)
    Iterator(OutputRow(key, updated))
  }

  override def close(): Unit = {}
}

Ключевые сущности:

  • StatefulProcessor[K, I, O] — абстрактный класс с lifecycle-методами. Наследуется один раз для каждой stateful-логики.
  • StatefulProcessorHandle — API для регистрации state-переменных и таймеров. Доступен через getHandle внутри процессора.
  • State-переменные — объявляются один раз в init() и кэшируются как поля класса. Каждая получает уникальное имя (используется в чекпоинте).
  • TimerValues — передаётся в handleInputRows, содержит текущее processing time и event time watermark.

Применение к DataFrame:

val result = events
  .groupByKey(_.userId)
  .transformWithState(
    new MyProcessor(),
    TimeMode.ProcessingTime(),
    OutputMode.Append()
  )

Заметьте: outputMode и timeMode — параметры самого оператора, а не глобального query, и попадают в init().

Lifecycle hooks StatefulProcessor

В отличие от v1, где была одна update-функция, v2 явно разделяет фазы жизненного цикла:

МетодКогда вызываетсяЗачем
init(outputMode, timeMode)Один раз при старте taskОбъявить state-переменные через getHandle
handleInitialState(key, initialState, timerValues)До первой обработки, если задан initial state batchPre-populate state из batch DataFrame
handleInputRows(key, rows, timerValues)Каждый micro-batch для ключа с новыми строкамиОсновная бизнес-логика
handleExpiredTimer(key, timerValues, expiredTimerInfo)Когда сработал зарегистрированный таймерОбработка timeout, эмиссия закрывающих событий
close()При завершении taskОсвобождение ресурсов
TIP

Почему два метода вместо одного. В mapGroupsWithState обработка новых событий и обработка timeout жили в одной функции, и приходилось писать if state.hasTimedOut: handle_timeout(); return. В v2 это два разных метода — логика разделена и тестируется независимо. Это особенно ценно, когда таймеров несколько и каждый делает что-то своё.

State-типы

Главное отличие v2 от v1 — несколько state-переменных на ключ и композитные типы, поддерживаемые state store нативно (без сериализации всего объекта при каждой операции).

ValueState[T] — одно значение

Эквивалент legacy GroupState[T], но без таймаут-семантики и с возможностью держать несколько штук на ключ.

private var lastSeen: ValueState[Timestamp] = _
private var sessionStart: ValueState[Timestamp] = _

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
  lastSeen = getHandle.getValueState("lastSeen", Encoders.TIMESTAMP)
  sessionStart = getHandle.getValueState("sessionStart", Encoders.TIMESTAMP)
}

Операции: get(), update(v), clear(), exists().

ListState[T] — упорядоченный список

Хранит коллекцию элементов под ключом. Главное преимущество: partial updates — добавление/удаление элементов без перечитывания и пересериализации всего списка.

private var eventBuffer: ListState[Event] = _

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
  eventBuffer = getHandle.getListState("buffer", Encoders.product[Event])
}

override def handleInputRows(
  key: String,
  inputRows: Iterator[Event],
  timerValues: TimerValues
): Iterator[OutputRow] = {
  // append без чтения всего списка -- O(1)
  inputRows.foreach(e => eventBuffer.appendValue(e))
  Iterator.empty
}

Операции: get() (итератор), appendValue(v), appendList(values), put(values), clear().

MapState[K, V] — словарь под ключом

Вложенный Map[K, V] под grouping key с point lookups: getValue(k) оптимизирован под одну запись, не требует загрузки всей мапы.

private var featureFlags: MapState[String, Boolean] = _

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
  featureFlags = getHandle.getMapState(
    "flags",
    Encoders.STRING,
    Encoders.scalaBoolean
  )
}

override def handleInputRows(
  key: String,
  inputRows: Iterator[FlagUpdate],
  timerValues: TimerValues
): Iterator[OutputRow] = {
  inputRows.foreach { upd =>
    featureFlags.updateValue(upd.flag, upd.value)
  }
  Iterator.empty
}

Операции: getValue(k), containsKey(k), updateValue(k, v), removeKey(k), keys(), values(), iterator(), clear().

INFO

Почему partial updates важны. В v1 чтобы добавить элемент в список из 10 000 элементов, нужно было прочитать весь список из state store, десериализовать его, добавить элемент и записать обратно. В v2 ListState.appendValue — одна запись в RocksDB column family. На больших коллекциях это уменьшает нагрузку на state store на порядки.

Таймеры (не state-переменная, но через тот же handle)

Таймеры — ещё один тип «состояния», но управляются отдельно через getHandle:

override def handleInputRows(
  key: String,
  inputRows: Iterator[Event],
  timerValues: TimerValues
): Iterator[OutputRow] = {
  // Зарегистрировать heartbeat через 5 минут (processing time)
  val now = timerValues.getCurrentProcessingTimeInMs()
  getHandle.registerTimer(now + 5 * 60 * 1000)

  // Зарегистрировать закрытие сессии по event time через 30 мин
  val watermark = timerValues.getCurrentWatermarkInMs()
  getHandle.registerTimer(watermark + 30 * 60 * 1000)

  Iterator.empty
}

API таймеров на handle:

  • registerTimer(expiryMs: Long) — зарегистрировать новый таймер.
  • deleteTimer(expiryMs: Long) — удалить ранее зарегистрированный.
  • listTimers() — получить все активные таймеры для текущего ключа.

Когда processing time >= expiryMs (или watermark >= expiryMs в event-time режиме), Spark вызовет handleExpiredTimer:

override def handleExpiredTimer(
  key: String,
  timerValues: TimerValues,
  expiredTimerInfo: ExpiredTimerInfo
): Iterator[OutputRow] = {
  val firedAt = expiredTimerInfo.getExpiryTimeInMs()
  // ... закрыть сессию, эмитнуть финальную запись и т.п.
  Iterator(OutputRow(key, "session_closed", firedAt))
}
TIP

Несколько таймеров на ключ — это новое. В mapGroupsWithState был ровно один timeout на ключ. В v2 их может быть произвольное количество, и каждый идентифицируется своим expiry timestamp. Это позволяет, например, отправлять ступенчатые алерты (через 5/15/60 минут бездействия) одной логикой.

TTL: декларативное удаление просроченного state

Каждая state-переменная может объявить time-to-live. При обновлении значения Spark автоматически записывает expiration time во вспомогательное хранилище и при следующих чтениях возвращает «нет значения», если TTL истёк. Никакого ручного state.hasTimedOut — логика прячется внутри state store.

import org.apache.spark.sql.streaming.TTLConfig
import java.time.Duration

private var sessionData: ValueState[SessionData] = _

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
  sessionData = getHandle.getValueState(
    "session",
    Encoders.product[SessionData],
    TTLConfig(ttlDuration = Duration.ofMinutes(30))
  )
}

После 30 минут без обновления значение автоматически исчезнет из state store на ближайшем micro-batch. Внутри Spark поддерживает отдельную TTL-структуру и периодически вызывает clearExpiredStateForAllKeys, чтобы физически удалять просроченные записи — иначе RocksDB рос бы бесконечно.

WARNING

TTL != timer. TTL только удаляет значение из state — он не даёт callback в код. Если нужно эмитнуть событие при истечении (например, «session closed»), регистрируйте registerTimer параллельно с TTL. TTL хорош для пассивной чистки, таймер — для активных reactions.

TTL поддерживается всеми типами: ValueState, ListState (per-list TTL) и MapState (per-entry TTL, начиная с 4.1).

PySpark: transformWithStateInPandas

В Python оператор называется transformWithStateInPandas (ещё есть row-based transformWithState, но pandas-вариант обычно быстрее за счёт векторизации). API структурно идентичен Scala-варианту.

import pandas as pd
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, StringType, LongType
from typing import Iterator

class FraudDetector(StatefulProcessor):
    def init(self, handle: StatefulProcessorHandle) -> None:
        # Объявляем три независимых state-переменных
        self.attempt_count = handle.getValueState(
            "attempt_count", "long"
        )
        self.recent_ips = handle.getListState(
            "recent_ips",
            StructType([StructField("ip", StringType())])
        )
        self.country_seen = handle.getMapState(
            "country_seen",
            StructType([StructField("country", StringType())]),
            "long"
        )

    def handleInputRows(
        self,
        key,
        rows: Iterator[pd.DataFrame],
        timer_values
    ) -> Iterator[pd.DataFrame]:
        count = self.attempt_count.get() if self.attempt_count.exists() else 0

        for batch in rows:
            for _, row in batch.iterrows():
                count += 1
                self.recent_ips.appendValue((row["ip"],))
                # счётчик стран
                prev = (
                    self.country_seen.getValue((row["country"],))
                    if self.country_seen.containsKey((row["country"],))
                    else 0
                )
                self.country_seen.updateValue((row["country"],), prev + 1)

        self.attempt_count.update(count)

        # Алерт, если 3+ failed login за batch
        if count >= 3:
            yield pd.DataFrame({
                "user_id": [key[0]],
                "alert": ["potential_fraud"],
                "attempts": [count]
            })

    def handleExpiredTimer(self, key, timer_values, expired_timer_info):
        # nothing for this example
        return iter([])

    def close(self) -> None:
        pass

output_schema = StructType([
    StructField("user_id", StringType()),
    StructField("alert", StringType()),
    StructField("attempts", LongType()),
])

alerts = (
    events
    .groupBy("user_id")
    .transformWithStateInPandas(
        statefulProcessor=FraudDetector(),
        outputStructType=output_schema,
        outputMode="Append",
        timeMode="ProcessingTime"
    )
)
TIP

Python-паритет наконец-то. В legacy-мире applyInPandasWithState (PySpark 3.4+) был ограниченным портом Scala-API, без полноценной TTL и таймеров. transformWithStateInPandas в Spark 4.0 обеспечивает функциональный паритет: те же ValueState/ListState/MapState, те же таймеры, те же TTL.

Полный пример: session window state machine

Соберём целиком: трекер пользовательских сессий с heartbeat (промежуточные эмиссии каждые 5 минут) и финальным закрытием по бездействию (30 минут).

import org.apache.spark.sql.{Encoders, SparkSession}
import org.apache.spark.sql.streaming._
import java.time.Duration
import java.sql.Timestamp

case class ClickEvent(userId: String, page: String, ts: Timestamp)
case class SessionUpdate(
  userId: String,
  startTs: Timestamp,
  lastTs: Timestamp,
  pageCount: Long,
  status: String
)

class SessionTracker
  extends StatefulProcessor[String, ClickEvent, SessionUpdate] {

  @transient private var sessionStart: ValueState[Timestamp] = _
  @transient private var lastSeen: ValueState[Timestamp] = _
  @transient private var pageBuffer: ListState[String] = _

  override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
    sessionStart = getHandle.getValueState(
      "sessionStart",
      Encoders.TIMESTAMP,
      TTLConfig(Duration.ofHours(2))   // safety net
    )
    lastSeen = getHandle.getValueState("lastSeen", Encoders.TIMESTAMP)
    pageBuffer = getHandle.getListState("pages", Encoders.STRING)
  }

  override def handleInputRows(
    key: String,
    inputRows: Iterator[ClickEvent],
    timerValues: TimerValues
  ): Iterator[SessionUpdate] = {
    val events = inputRows.toList
    if (events.isEmpty) return Iterator.empty

    // 1. Старт сессии при первом событии
    if (!sessionStart.exists()) {
      sessionStart.update(events.head.ts)
    }

    // 2. Аккумуляция страниц через ListState
    events.foreach(e => pageBuffer.appendValue(e.page))

    // 3. Обновление last seen
    val newest = events.maxBy(_.ts.getTime).ts
    lastSeen.update(newest)

    // 4. Перезапуск таймеров: heartbeat и close
    getHandle.listTimers().forEach(t => getHandle.deleteTimer(t))
    val nowMs = timerValues.getCurrentProcessingTimeInMs()
    getHandle.registerTimer(nowMs + 5 * 60 * 1000)   // heartbeat
    getHandle.registerTimer(nowMs + 30 * 60 * 1000)  // close

    val pageCount = pageBuffer.get().asScala.size
    Iterator(SessionUpdate(
      key, sessionStart.get(), newest, pageCount, "active"
    ))
  }

  override def handleExpiredTimer(
    key: String,
    timerValues: TimerValues,
    expiredTimerInfo: ExpiredTimerInfo
  ): Iterator[SessionUpdate] = {
    val nowMs = expiredTimerInfo.getExpiryTimeInMs()
    val sessionAgeMs = nowMs - lastSeen.get().getTime
    val pages = pageBuffer.get().asScala.toList

    if (sessionAgeMs >= 30 * 60 * 1000) {
      // Финал: закрытие сессии и очистка
      val update = SessionUpdate(
        key, sessionStart.get(), lastSeen.get(), pages.size.toLong, "closed"
      )
      sessionStart.clear()
      lastSeen.clear()
      pageBuffer.clear()
      Iterator(update)
    } else {
      // Heartbeat
      Iterator(SessionUpdate(
        key, sessionStart.get(), lastSeen.get(), pages.size.toLong, "heartbeat"
      ))
    }
  }

  override def close(): Unit = {}
}

val spark = SparkSession.builder
  .appName("SessionTrackerV2")
  .config(
    "spark.sql.streaming.stateStore.providerClass",
    "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider"
  )
  .getOrCreate()

import spark.implicits._

val events = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "kafka:9092")
  .option("subscribe", "clickstream")
  .load()
  .selectExpr("CAST(value AS STRING) AS json")
  .select(from_json($"json", Encoders.product[ClickEvent].schema).as("e"))
  .select($"e.*").as[ClickEvent]

val sessions = events
  .groupByKey(_.userId)
  .transformWithState(
    new SessionTracker(),
    TimeMode.ProcessingTime(),
    OutputMode.Update()
  )

sessions.writeStream
  .outputMode("update")
  .format("delta")
  .option("checkpointLocation", "/checkpoints/sessions-v2")
  .start("/data/gold/sessions_v2")

Обратите внимание: state хранится в трёх разных переменных (start/lastSeen/pages), не в одном «упакованном» объекте. Каждая переменная читается и обновляется независимо, что и составляет основное архитектурное преимущество v2.

State store backends и Coordinated Commits

transformWithState поддерживается только на RocksDB state store provider. HDFS-провайдер (in-memory) для v2 не работает — v2 опирается на column families и range scans RocksDB, которых у HDFS-провайдера нет.

spark.conf.set(
  "spark.sql.streaming.stateStore.providerClass",
  "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider"
)

Внутри RocksDB каждая state-переменная (ValueState, ListState, MapState) живёт в своей column family — логически изолированной части DB с собственным compaction. Это и обеспечивает partial updates без чтения смежных переменных.

Дополнительно transformWithState интегрирован с Coordinated Commits — механизмом, появившимся в Spark 4.0 для атомарной фиксации checkpoint state и output sink. На Databricks он включается автоматически, в OSS Spark включается через spark.sql.streaming.stateStore.coordinatedCommitsEnabled = true.

WARNING

Только RocksDB. Если вы запустите transformWithState с дефолтным HDFSBackedStateStoreProvider, query упадёт ещё на стадии планирования. RocksDB — обязательное требование.

Migration: mapGroupsWithState -> transformWithState

Прямой автоматической миграции не существует — API слишком разные. Но рецепт пошагового перехода такой:

  1. Замените функцию на класс. Тело (key, events, state) => ... оборачивается в StatefulProcessor с init/handleInputRows.
  2. Распакуйте состояние. Каждое поле вашего бывшего case class вынесите в отдельную ValueState. Если внутри был список или мапа — замените их на ListState/MapState ради partial updates.
  3. Замените timeout на TTL или таймер. state.setTimeoutDuration("30 min") -> либо TTLConfig(Duration.ofMinutes(30)) (если нужно только удалить), либо registerTimer(now + 30*60*1000) + handleExpiredTimer (если нужно эмитнуть событие).
  4. Разделите логику timeout и input. Уберите if state.hasTimedOut: ... return — его место теперь в handleExpiredTimer.
  5. Новый чекпоинт. Запустите v2-query с новым checkpointLocation: state schema несовместима со старым. Параллельно гоняйте обе версии до подтверждения корректности — legacy-чекпоинт продолжать работать на v1-операторе.
WARNING

Несовместимость чекпоинтов. transformWithState использует другую структуру state store (column families, отдельные TTL-store). Восстановление из старого mapGroupsWithState-чекпоинта невозможно. Планируйте dual-run или backfill.

State schema evolution

Отдельный плюс v2: схема state может эволюционировать, если включён Avro-кодинг.

spark.conf.set(
  "spark.sql.streaming.stateStore.encodingFormat",
  "avro"
)

С Avro можно добавлять новые опциональные поля в ValueState[T], расширять числовые типы (Int -> Long) и удалять неиспользуемые поля — при рестарте Spark прочитает старые записи через Avro schema resolution. С дефолтным unsafe-кодингом любое изменение полей — breaking change.

Сравнение API

СвойствоmapGroupsWithStateflatMapGroupsWithStatetransformWithState (v2)
ВерсияSpark 2.2+Spark 2.2+Spark 4.0 GA
Output cardinalityРовно 1 строка на вызов0/1/N строк0/1/N строк
State-переменных на ключ1 (один GroupState[S])1Несколько
Композитные state-типыНет (только в одном объекте)НетValueState, ListState, MapState
Partial updatesНетНетДа (per-element для List/Map)
Timer’ов на ключ11Произвольно много
TTLРучной (hasTimedOut)РучнойДекларативный, per-state
Lifecycle hooksОдин callbackОдин callbackinit/handleInputRows/handleExpiredTimer/handleInitialState/close
Initial stateЧерез mapGroupsWithStateInitialStateАналогичноНативно через handleInitialState
State storeHDFS или RocksDBHDFS или RocksDBТолько RocksDB
Schema evolutionНетНетДа, через Avro encoding
Python parityТолько через applyInPandasWithStateАналогичноtransformWithStateInPandas (полный паритет)
StatusMaintenanceMaintenanceRecommended
INFO

Когда оставаться на legacy. Если у вас уже работает mapGroupsWithState в production и логика простая (один state, один timeout), переезжать ради переезда не нужно: v1 продолжит поддерживаться. Но любой новый stateful pipeline в Spark 4.0+ стоит сразу писать на transformWithState — API мощнее, и Databricks/OSS Spark будут вкладываться именно в него.

Что дальше в Spark 4.x

Roadmap по Arbitrary State API v2 (по тикетам Apache JIRA и анонсам Databricks):

  • Spark 4.1 — state schema evolution через Avro (готово), chaining stateful operators (transformWithState -> groupBy -> transformWithState), расширенный state data source для batch read и предварительная поддержка batch write.
  • Spark 4.2 (планируется) — streaming-side write через state data source (миграция state’а без рестарта), улучшения наблюдаемости (per-state-variable метрики), batch state preview.
  • ДолгосрочноtransformWithState рассматривается как фундамент для будущих stateful-операторов, включая встроенные сессии, дедупликацию и изменения в SQL-уровне.

Knowledge Check

Проверка знанийKnowledge check
Почему transformWithState нельзя запустить с HDFS-backed state store, который работал для mapGroupsWithState?
ОтветAnswer
transformWithState опирается на возможности RocksDB, которых нет у in-memory HDFS-провайдера: column families (изоляция переменных state), range scans (для эффективных операций на ListState/MapState и для TTL-store) и merge operators. Каждая ValueState/ListState/MapState физически живёт в своей column family и взаимодействует с RocksDB API напрямую. HDFS-провайдер устроен как plain Map[Key, Value] поверх памяти и фиксированных снапшотов, и встроить туда composite state и declarative TTL без полной переписки невозможно. Поэтому для transformWithState обязательна конфигурация spark.sql.streaming.stateStore.providerClass = RocksDBStateStoreProvider.
Проверка знанийKnowledge check
В чём разница между TTL и таймером в transformWithState? Когда использовать каждый?
ОтветAnswer
TTL -- это пассивная политика чистки: Spark автоматически удаляет state-значение после истечения duration, никакого callback в коде не вызывается. Таймер -- это активная регистрация будущего события: Spark в нужный момент вызывает handleExpiredTimer и даёт коду возможность эмитнуть выходную строку, обновить state, перерегистрировать таймер и т.д. TTL подходит, когда достаточно просто избавиться от мёртвого state (deduplication windows, протухшие feature flags, кэши). Таймер нужен, когда требуется реакция: эмитнуть финальную запись о сессии, отправить heartbeat, послать алерт. На практике их часто комбинируют: TTL ставят как safety net (state не утечёт, даже если логика таймера сломалась), а таймер -- для основной business logic.
Проверка знанийKnowledge check
Какие архитектурные ограничения mapGroupsWithState решает transformWithState за счёт нескольких state-переменных и partial updates?
ОтветAnswer
В mapGroupsWithState вся информация о ключе паковалась в один объект GroupState[S], и любая операция (даже добавление одного элемента в список внутри S) требовала прочитать весь объект, десериализовать, изменить и записать целиком. При больших структурах (списках событий, мапах сессий) это давало O(N) IO на каждое обновление и пробивало пропускную способность state store. transformWithState решает это двумя способами: (1) разные логические аспекты состояния -- lastSeen, sessionStart, pageBuffer -- живут в разных state-переменных и читаются независимо, поэтому обновление одной не трогает другие; (2) ListState и MapState поддерживают partial updates -- appendValue/updateValue работают за O(1) к одному элементу, без чтения всей коллекции. В сумме это снижает нагрузку на RocksDB на порядки в сценариях с длинными буферами и большими словарями.

Что дальше?

В следующем уроке мы перейдём к CDC и Debezium — разберём, как потреблять Change Data Capture-события из Debezium через Kafka и применять transformWithState для построения консистентных Delta Lake snapshot’ов из CDC-потока.

Debezium → PySpark Structured Streaming

Закончили урок?

Отметьте его как пройденный, чтобы отслеживать свой прогресс

Войдите чтобы оценить урок

Прогресс модуля
0 из 8