skip to content

How do you add a field to a POJO stored in Flink ValueState without losing the state?

level: seniorimportance: should knowfreq 40%

answer

  1. The upgrade artifact is a savepoint
  2. Only two type families qualify
  3. Added fields get the Java default
  4. Renaming or repackaging the class breaks it
  5. Without stable uids nothing matches back

basics

~20 s

Take a savepoint, add the field to the POJO, and restore from that savepoint. Flink evolves POJO and Avro state schemas automatically: added fields get their Java default value, removed fields are dropped, but declared field types and the class name cannot change.

solid answer

~50 s

The procedure is three steps: take a savepoint of the running job, change the state type in the code, restore from the savepoint. On first access to that state, Flink compares the new serializer's schema with the persisted one and, if they differ, reads the old bytes with the previous serializer and writes them back with the new one. It works because Flink's own type serialization framework produced the serializer — so the state descriptor must declare the type (`new ValueStateDescriptor<>("profile", Profile.class)`), not a hand-written `TypeSerializer`. Schema evolution is supported only for **POJO** and **Avro** types. The POJO rules: fields may be added (initialised to the Java default) or removed, but a declared field's type cannot change and the class name, including its package, must stay the same. Java records follow the same rules. Set stable operator `uid`s or the state will not be matched back at all.

code

java · 14 lines
java
// Before: state type as deployed
public class Profile {
    public String userId;
    public long visits;
}

// After: a field added -- allowed.
// tier restores as null for every pre-existing entry,
// NOT as whatever a constructor would set.
public class Profile {
    public String userId;
    public long visits;
    public String tier;
}

go deeper

for a junior

Recall the three-step upgrade: savepoint, change the type, restore from the savepoint. Knowing that Flink can do this at all for POJOs is the level-appropriate answer.

for a middle

Explain the mechanics: Flink compares serializer schemas at first access and rewrites with the new serializer, it only works for POJO and Avro, and added fields get Java defaults rather than constructor values.

for a senior

Show that you have upgraded a real job — stable operator uids as the precondition, Kryo contagion through generic fields as the trap, and the State Processor API as the escape hatch for type changes Flink will not perform.

for a principal

Own the compatibility policy: which serialization format state types must use, how schema changes are reviewed before deployment, whether a savepoint-rewrite step belongs in the release pipeline, and the rollback plan when a migration goes wrong.

## The problem Streaming jobs run for months, and the data classes they keep in state change with the product. If you add a field to the POJO stored in a `ValueState` and simply redeploy, one of three things happens: the restore fails with a state-migration error, the state is silently abandoned, or Flink migrates it for you. Which one depends entirely on the serializer behind the state. ## The mechanism Whether state schema can evolve is decided by the serializer that reads and writes the persisted bytes. Flink handles this transparently for serializers generated by its own type serialization framework — that is, when the `StateDescriptor` you register declares a type and lets Flink infer the serializer: ```java ValueStateDescriptor<Profile> descriptor = new ValueStateDescriptor<>("profile", Profile.class); ``` If you instead pass an explicit `TypeSerializer` or `TypeInformation`, you own compatibility yourself, and you have to implement the serializer-snapshot machinery to declare what is compatible. The migration itself happens at first access to a given state after restore, independently for each registered state: Flink checks whether the new serializer's schema differs from the previous one; if it does, it reads the persisted objects with the *previous* serializer and writes them back with the *new* one. ## The procedure 1. Take a savepoint of the running job. 2. Update the state types in the application — add the field, regenerate the Avro schema, whatever the change is. 3. Restore the job from the savepoint. Savepoint, not checkpoint, is the documented path — a savepoint is the artifact intended for job upgrades. ## What is supported Schema evolution covers exactly two families of types: **POJO** and **Avro**. If state schema evolution matters to you, that is a strong reason to use one of them for anything you store. **POJO types** evolve under four rules: 1. Fields can be removed. The removed field's previous value is dropped in future checkpoints and savepoints. 2. New fields can be added. Each is initialised to the Java default for its type — `0`, `false`, `null` — not to any value your constructor would produce. If "absent" and "zero" mean different things in your logic, make the new field a boxed type so absent reads as `null`, or backfill it explicitly on first access. 3. Declared field types cannot change. Widening `int` to `long` is a schema change Flink will not perform. 4. The class name cannot change, including the namespace. Moving the class to another package is a breaking change. In Flink 2.3, Java records are treated as POJO types and follow the same rules. **Avro types** are fully supported as long as Avro's own rules for schema resolution consider the change compatible — which gives you defaults for added fields, so an added Avro field can carry a meaningful default rather than a Java zero. The limitation is that Avro-generated classes used as the state type cannot be relocated or given a different namespace on restore. ## What is not supported **Keys cannot evolve.** The structure of a key is off limits, because migrating it can be non-deterministic: drop a field from a POJO used as a key and two previously distinct keys may become identical, with no rule for merging their values. On top of that, `EmbeddedRocksDBStateBackend` compares keys by their serialised bytes rather than through `hashCode()`, so any change to a key's object structure changes the mapping. **Kryo cannot evolve.** When a type falls back to Kryo, Flink has no way to verify whether a change is compatible, so evolution is simply unsupported. The trap is that this is contagious through containers: when a field falls back to Kryo, its contents go through Kryo too, so a well-behaved POJO inside it cannot be evolved. What falls back changed in Flink 2.0, which added built-in serializers for collections: - A field declared through the `List`, `Map` or `Set` interface with concrete type arguments — `List<SomeOtherPojo>` — gets Flink's own collection serializer, and `SomeOtherPojo` keeps the POJO rules. The schema-evolution documentation page still gives exactly this field as its Kryo example, but the 2.3 type extractor no longer routes it to Kryo. - A field declared as a concrete implementation class (`ArrayList<SomeOtherPojo>`), a raw `List`, or a collection of a type variable still falls back to Kryo, contents and all. Auditing which of your state types actually resolve to Kryo — Flink logs generic-type fallbacks — is worth doing before you promise anyone a smooth upgrade path. ## The prerequisite people forget: operator uids None of this matters if Flink cannot find the state in the first place. Savepoints map stored state to operators by uid. By default uids are generated by traversing the JobGraph and hashing operator properties, which is convenient and extremely fragile: inserting an operator, or swapping one, changes the generated uid and the old state no longer matches. Call `uid(String)` on every stateful operator from the very first deployment. Retrofitting uids onto a running job is painful, because the auto-generated hashes are what the existing savepoint recorded. ## When the rules do not cover you Two escape hatches. Write a custom `TypeSerializer` with a serializer snapshot that declares the compatibility you know to be safe — appropriate when you understand the byte format and the framework's conservatism is the only obstacle. Or use the **State Processor API**: read the savepoint offline with `SavepointReader`, transform the entries in a batch job, and write a new savepoint with `SavepointWriter`. That is the tool for changes Flink refuses to do automatically — changing a field's type, changing the key, splitting one state into two, or moving from `ValueState<Map<K,V>>` to `MapState<K,V>`. ## A related migration that used to hurt Turning TTL on or off for an existing state used to fail the restore with a `StateMigrationException`, because TTL-enabled state uses a different serialisation format. From Flink 2.2.0 that migration is seamless in both directions across the heap and RocksDB backends: restoring non-TTL state under a TTL-enabled descriptor treats the existing entries as not expired, so each one starts expiring only after its next access or update, and going the other way simply ignores the TTL metadata. Changing TTL *parameters* is still not always compatible. ## Answering well Give the three-step procedure, name POJO and Avro as the supported families, state the four POJO rules crisply, and then show operating judgment: mention uids as the precondition, Kryo contagion as the trap, and the State Processor API as the escape hatch for everything else.

  • Why can't the type of an existing POJO field be changed?
    Flink's POJO serializer persists each field with the serializer for its declared type, and it will not attempt a value-level conversion between two different types — even a widening one such as `int` to `long`. The framework can only add fields with their Java defaults and drop fields it no longer needs. To change a field's type you must rewrite the savepoint yourself with the State Processor API, or migrate through a new field and drop the old one in a later release.
  • What happens if a POJO stored in state contains a List of another POJO?
    Since Flink 2.0 it depends on the declaration. Declared through the interface with a concrete type argument — `List<SomeOtherPojo>` — it gets Flink's built-in list serializer, so `SomeOtherPojo` keeps the POJO evolution rules. Declared as a concrete class such as `ArrayList<SomeOtherPojo>`, raw, or over a type variable, it falls back to Kryo, which cannot evolve, and the contained type is frozen.
  • Why must every stateful operator have an explicit uid?
    Savepoints map stored state back to operators by uid. Without `uid(...)`, Flink generates uids by hashing operator properties from the JobGraph, so almost any topology edit — inserting an operator, reordering, swapping one — changes the hash and the operator restores with empty state. Setting stable uids from the first deployment is a production-readiness requirement, not an optimisation; retrofitting them onto a live job is painful.
  • How do you make a change Flink refuses to migrate automatically?
    Use the State Processor API. Load the savepoint with `SavepointReader` in a batch job, read out the keyed and operator state, transform it however you need — change a field type, re-key it, split one state into two — then write a fresh savepoint with `SavepointWriter` and start the new job from it. It is the sanctioned escape hatch for anything outside the POJO and Avro rules.

saying these in an interview costs you the question

  • Says any state type can evolve if you just redeploy
  • Expects an added POJO field to get its constructor value
  • Thinks renaming or repackaging the class is harmless
  • Believes Kryo-serialized types evolve like POJOs
  • Tries to evolve the schema of the state key

context