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 вместо настоящего
mapGroupsWithState—applyInPandasWithState, который привязан к 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 и десятках связанных тикетов) решает эти проблемы за счёт перехода от «функции обновления» к объекту-процессору с явным жизненным циклом.
Версии. Оператор 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 batch | Pre-populate state из batch DataFrame |
handleInputRows(key, rows, timerValues) | Каждый micro-batch для ключа с новыми строками | Основная бизнес-логика |
handleExpiredTimer(key, timerValues, expiredTimerInfo) | Когда сработал зарегистрированный таймер | Обработка timeout, эмиссия закрывающих событий |
close() | При завершении task | Освобождение ресурсов |
Почему два метода вместо одного. В 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().
Почему 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))
}
Несколько таймеров на ключ — это новое. В 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 рос бы бесконечно.
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"
)
)
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.
Только RocksDB. Если вы запустите transformWithState с дефолтным HDFSBackedStateStoreProvider, query упадёт ещё на стадии планирования. RocksDB — обязательное требование.
Migration: mapGroupsWithState -> transformWithState
Прямой автоматической миграции не существует — API слишком разные. Но рецепт пошагового перехода такой:
- Замените функцию на класс. Тело
(key, events, state) => ...оборачивается вStatefulProcessorсinit/handleInputRows. - Распакуйте состояние. Каждое поле вашего бывшего case class вынесите в отдельную
ValueState. Если внутри был список или мапа — замените их наListState/MapStateради partial updates. - Замените timeout на TTL или таймер.
state.setTimeoutDuration("30 min")-> либоTTLConfig(Duration.ofMinutes(30))(если нужно только удалить), либоregisterTimer(now + 30*60*1000)+handleExpiredTimer(если нужно эмитнуть событие). - Разделите логику timeout и input. Уберите
if state.hasTimedOut: ... return— его место теперь вhandleExpiredTimer. - Новый чекпоинт. Запустите v2-query с новым checkpointLocation: state schema несовместима со старым. Параллельно гоняйте обе версии до подтверждения корректности — legacy-чекпоинт продолжать работать на v1-операторе.
Несовместимость чекпоинтов. 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
| Свойство | mapGroupsWithState | flatMapGroupsWithState | transformWithState (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’ов на ключ | 1 | 1 | Произвольно много |
| TTL | Ручной (hasTimedOut) | Ручной | Декларативный, per-state |
| Lifecycle hooks | Один callback | Один callback | init/handleInputRows/handleExpiredTimer/handleInitialState/close |
| Initial state | Через mapGroupsWithStateInitialState | Аналогично | Нативно через handleInitialState |
| State store | HDFS или RocksDB | HDFS или RocksDB | Только RocksDB |
| Schema evolution | Нет | Нет | Да, через Avro encoding |
| Python parity | Только через applyInPandasWithState | Аналогично | transformWithStateInPandas (полный паритет) |
| Status | Maintenance | Maintenance | Recommended |
Когда оставаться на 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
Что дальше?
В следующем уроке мы перейдём к CDC и Debezium — разберём, как потреблять Change Data Capture-события из Debezium через Kafka и применять transformWithState для построения консистентных Delta Lake snapshot’ов из CDC-потока.