CertSafari
    Databricks Certified Associate Developer for Apache Spark· Lessons

    Domain 3 · Lesson 19/32

    Stateful operators and StateStores with transformWithState

    Create and invoke user-defined functions with or without stateful operators, including StateStores.

    12 min read
    3.12% of exam
    3 sources
    Published 3 Oct 2026
    Docs as of 30 Sep 2026

    What you will be able to do

    • Distinguish stateful from stateless streaming operations, and configure the RocksDB state store provider
    • Define a StatefulProcessor whose state variables use ValueState, ListState, or MapState, and predict how TTL behaves for each
    • Choose between handleInputRows, handleExpiredTimer, and handleInitialState
    • Invoke a processor with groupBy().transformWithStateInPandas() or transformWithState, and recognize the legacy operators it replaces

    1.Stateful queries and the state store

    A stateless streaming query only tracks which source rows it has already processed. A stateful query also keeps intermediate state information and updates it incrementally with each micro-batch. That covers streaming aggregations, distinct, dropDuplicates, stream-stream joins, and custom stateful applications. The state lives in a state store. Its provider is set by spark.sql.streaming.stateStore.providerClass. RocksDB is the default provider in Databricks Runtime 17.3 and above, and on earlier runtimes you have to configure it yourself. Databricks also recommends RocksDB with changelog checkpointing for stateful streams.

    Enabling the RocksDB state store provider for the session (needed below DBR 17.3)python
    spark.conf.set("spark.sql.streaming.stateStore.providerClass", "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")

    The checkpoint fixes the shuffle partition count in the same way. Changing spark.sql.shuffle.partitions has no effect on a query that already has a checkpoint. Before you write custom state logic at all, note what Databricks advises: for aggregations, deduplication, and streaming joins, use the built-in operators. Custom stateful processors are for logic those operators can't express.

    Checkpoint 1 of 7· Check yourself

    You need to drop duplicate events from a stream. Which approach does Databricks recommend?

    Sources12

    2.Defining a StatefulProcessor and its state variables

    With transformWithState, you write a class that extends StatefulProcessor. Spark passes a StatefulProcessorHandle to your init method, and you use that handle to create state variables with getValueState, getListState, or getMapState. Each state variable needs a unique name and a schema. In Python the schema is mandatory, and in Scala you can pass an Encoder instead. You can optionally add a TTL in milliseconds. A MapState needs separate schemas for its keys and its values. One processor can hold several state variables. The grouping key comes from the groupBy before the operator, and Spark keeps every state variable separately for each grouping key.

    A Python processor with a MapState keyed by session_id, grouped by user_idpython
    class SessionTracker(StatefulProcessor):
      def init(self, handle: StatefulProcessorHandle) -> None:
        self.sessions = handle.getMapState("sessions", "session_id string", "count long")
    
      def handleInputRows(self, key, rows: Iterator[Row], timerValues) -> Iterator[Row]:
        for row in rows:
          session_key = (row["session_id"],)  # session_id is the MapState key
          count = self.sessions.getValue(session_key)[0] if self.sessions.containsKey(session_key) else 0
          new_count = count + 1
          self.sessions.updateValue(session_key, (new_count,))
        yield from []
    
      def close(self) -> None:
        pass
    
    df.groupBy("user_id").transformWithState(SessionTracker(), ...) # user_id is the grouping key
    The three state types in transformWithState
    State typeStores per grouping keyHow you change itTTL granularity
    ValueStateOne value (can be a struct or tuple)Replace the entire valueResets when the value updates
    ListStateA list of valuesAppend an item, append a list, or overwrite with putEach list value has its own TTL; only put resets it
    MapStateDistinct keys, each mapped to a valueUpdate or remove a key; list keys, values, or iterate pairsEach key-value pair has its own TTL

    Checkpoint 2 of 7· Exam question

    Which `timeoutConf` value should be passed to `applyInPandasWithState` so a group's state expires after a fixed amount of wall-clock time has elapsed since it last received data, independent of any event-time watermark in the incoming records?

    Checkpoint 3 of 7· Match them up

    Match each requirement to the state type that fits it

    Tap a term, then the definition that fits it.

    Sources2

    3.handleInputRows, handleExpiredTimer and handleInitialState

    Your processing logic goes in its methods. handleInputRows runs for each grouping key that has rows in the micro-batch. It gets those rows as an iterator, can read and write state, and yields output rows. handleExpiredTimer runs when a timer fires, whether or not the key received new rows. Timers let you go beyond simple state eviction, including emitting rows. handleInitialState is optional and pre-populates state before any input rows arrive.

    What each handler can do (from the Databricks comparison)
    BehaviorhandleInputRowshandleExpiredTimer
    Get, put, update, or clear state valuesYesYes
    Create or delete a timerYesYes
    Emit rowsYesYes
    Iterate over rows in the current micro batchYesNo
    Trigger logic based on elapsed timeNoYes

    Checkpoint 4 of 7· Check yourself

    Which handler can iterate over the rows of the current micro-batch?

    Sources2

    4.Invoking the processor on a streaming DataFrame

    PySpark offers two operators. transformWithState is row-based, and transformWithStateInPandas passes pandas DataFrames to your handlers. Scala supports only the row-based API. In Python you call either operator on df.groupBy(key) and supply the processor instance, an output schema, an output mode, and a time mode. The SCD type 1 example below keeps only the latest location per user in a ValueState.

    Checkpoint 5 of 7· Fill the gap

    Which handle method creates the single-value state used by this SCD type 1 processor?

            # Initialize the state to store the latest location for each user
            self.latest_location = handle. ? ("latestLocation", value_state_schema)
    Invoking a processor with groupBy().transformWithStateInPandas()python
    q = (
        df.groupBy("user")
        .transformWithStateInPandas(
            statefulProcessor=SCDType1StatefulProcessor(),
            outputStructType=output_schema,
            outputMode="Update",
            timeMode="None",
        )
        .writeStream.format("memory")
        .queryName("scd1_output")
        .option("checkpointLocation", f"/tmp/checkpoint_{uuid.uuid4()}")
        .trigger(availableNow=True)
        .start()
    )

    Checkpoint 6 of 7· Exam question

    A data engineer at a logistics company has a plain Python function: ```python def normalize_country(code: str) -> str: return code.strip().upper() ``` They need to apply it to the `country` column of a batch DataFrame `df` using the DataFrame API's `.withColumn()`. Which code correctly registers and applies the function?

    Sources2

    5.Legacy arbitrary stateful operators

    Before transformWithState, custom state ran through mapGroupsWithState and flatMapGroupsWithState, plus applyInPandasWithState in Python. In these operators the state update function receives the previous state as a GroupState object. Both Scala operators can accept a user-defined initial state, so a stream started without a valid checkpoint doesn't have to reprocess data. Databricks now recommends transformWithState for arbitrary state transformations. It is also the custom operator that works in pipelines that chain several stateful operators, in DBR 16.2 and later. Those pipelines don't support the legacy operators and allow only append output mode.

    Legacy flatMapGroupsWithState signature: a function over (key, values, GroupState)scala
    def flatMapGroupsWithState[S: Encoder, U: Encoder](
        outputMode: OutputMode,
        timeoutConf: GroupStateTimeout,
        initialState: KeyValueGroupedDataset[K, S])(
        func: (K, Iterator[V], GroupState[S]) => Iterator[U])

    Checkpoint 7 of 7· Check yourself

    A pipeline feeds a windowed aggregation into a custom stateful step. Which custom operator works in this multi-operator chain?

    Sources213

    Exam traps

    Each one states something that sounds right. Open it to see what is actually true.

    1. 1.Appending items to a ListState refreshes the TTL of the list.Why is that wrong?

      Each list value has its own TTL, and only a put operation, which overwrites the whole list, resets it.

      Covered in Defining a StatefulProcessor and its state variables

    2. 2.You can update a single field of a struct held in ValueState.Why is that wrong?

      A ValueState holds a single value per key, and your logic has to replace the entire value.

      Covered in Defining a StatefulProcessor and its state variables

    3. 3.Raising spark.sql.shuffle.partitions and restarting a checkpointed stateful query re-partitions its state.Why is that wrong?

      The partition count is fixed when the checkpoint is created. A query that already has a checkpoint keeps its original count.

      Covered in Stateful queries and the state store

    Practise it for real

    Run the Databricks SCD type 1 example and confirm that the state store keeps only the latest location per user.

    1. 1.Set spark.sql.streaming.stateStore.providerClass to the RocksDBStateStoreProvider class (required below DBR 17.3).

      Why: transformWithState needs RocksDB as the state store provider on older runtimes.

      You should see: The config call returns without error.

    2. 2.Define SCDType1StatefulProcessor, with init calling handle.getValueState("latestLocation", value_state_schema).

      Why: A ValueState holds exactly one row per user, which is what SCD type 1 needs.

      You should see: The class definition compiles in the notebook.

    3. 3.Create main.stateful_examples, seed scd1_source with rows (u1,1,NYC), (u1,3,SF), (u1,2,LA), (u2,5,London), and read it with spark.readStream.table.

      Why: This gives you a small streaming source with out-of-order times for u1.

      You should see: df is a streaming DataFrame.

    4. 4.Run df.groupBy("user").transformWithStateInPandas(...) writing to the memory sink scd1_output with trigger(availableNow=True), then call q.awaitTermination().

      Why: This invokes the processor once per grouping key and saves state under the checkpoint location.

      You should see: The query processes the available data and stops.

    5. 5.Query SELECT user, time, location FROM scd1_output ORDER BY user.

      Why: The output shows what the processor emitted from state.

      You should see: u1 -> SF (time 3) and u2 -> London (time 5).

    Stuck? Get a nudge

    If u1 shows LA, check that handleInputRows compares against the maximum time, not the last row it saw.

    Sources

    Every claim above is drawn from one of these pages, quoted as it was written on the date shown.

    1. 1.
      “Stateful operations include streaming aggregation, distinct, dropDuplicates, stream-stream joins, and custom stateful applications.”
      ↩︎ Stateful queries and the state store
      “Legacy custom stateful operators (FlatMapGroupWithState and applyInPandasWithState) are not supported.”
      ↩︎ Legacy arbitrary stateful operators
      “Changing spark.sql.shuffle.partitions has no effect on a streaming query that already has a checkpoint”
      ↩︎ Exam trap 3
      “The state management scheme can't be changed between query restarts.”
      ↩︎ Prediction
      “In Databricks Runtime 16.2 or later, you can use transformWithState in workloads with multiple stateful operators.”
      ↩︎ Checkpoint
    2. 2.
      “RocksDB is the default state store provider in Databricks Runtime 17.3 and above.”
      ↩︎ Stateful queries and the state store
      “Spark passes a StatefulProcessorHandle to the init method of your StatefulProcessor.”
      ↩︎ Defining a StatefulProcessor and its state variables
      “In Python, you must specify the schema.”
      ↩︎ Defining a StatefulProcessor and its state variables
      “Optionally, implement handleInitialState to pre-populate state before your application processes any input rows.”
      ↩︎ handleInputRows, handleExpiredTimer and handleInitialState
      “Timers allow you to define custom logic beyond state eviction, including emitting rows.”
      ↩︎ handleInputRows, handleExpiredTimer and handleInitialState
      “Scala supports only the row-based transformWithState API.”
      ↩︎ Invoking the processor on a streaming DataFrame
      “Databricks recommends using transformWithState instead of legacy operators, such as flatMapGroupsWithState and mapGroupsWithState, for arbitrary state transformations.”
      ↩︎ Legacy arbitrary stateful operators
      “To reset time-to-live, you must use a put operation.”
      ↩︎ Exam trap 1
      “For ValueState, you must implement logic to replace the entire value.”
      ↩︎ Exam trap 2
      “For stateful operations such as aggregations, deduplication, and streaming joins, Databricks recommends using built-in Structured Streaming operators instead of custom logic.”
      ↩︎ Checkpoint
      “If you process a source key for ValueState without updating the stored ValueState, the time-to-live isn't reset.”
      ↩︎ Prediction
      “transformWithState supports three state types: ValueState, ListState, and MapState.”
      ↩︎ Checkpoint
      “Implement handleExpiredTimer to run time-based logic regardless of whether the grouping key receives new rows in a micro-batch.”
      ↩︎ Checkpoint
      “transformWithStateInPandas is not supported in real-time mode. Instead, use transformWithState.”
      ↩︎ Prediction
    3. 3.
      “The state update function takes the previous state as input using an object of type GroupState.”
      ↩︎ Legacy arbitrary stateful operators

    Ready to test yourself?

    Practise Databricks Certified Associate Developer for Apache Spark in quiz mode.

    Spotted a mistake, or was something unclear? Tell us.