Thanks Shengkai and Leonard for the reviews and the +1s (and everyone for
the earlier discussion)!

If there are no further concerns, I'll start a [VOTE] thread in a few days.

Best,
Weiqing

On Mon, Jul 13, 2026 at 10:39 PM Leonard Xu <[email protected]> wrote:

> Hi Weiqing,
>
>
>
> Thanks for driving this FLIP, +1 for the proposal.
>
> I think this is a very useful feature for many production SQL jobs. In
> practice, it is common for input schemas to evolve as business requirements
> change, especially for CDC / dimension-table scenarios where adding
> nullable fields or evolving nested ROW structures should not necessarily
> require users to drop state or reprocess all historical data.
>
> The current scoped and opt-in design looks reasonable to me. It keeps the
> migration behavior safe by default, while still giving users a path to
> handle backward-compatible RowData schema changes when they understand the
> state mapping of their job.
>
> Overall, I see this as an important step towards better schema evolution
> support for Flink SQL jobs.
>
>
>
> Best,
>
> Leonard
>
>
> > 2026 7月 14 10:42 上午,Shengkai Fang <[email protected]> 写道:
> >
> > +1 for the proposal.
> >
> > Best,
> > Shengkai
> >
> > Weiqing Yang <[email protected]> 于2026年7月9日周四 14:14写道:
> >
> >> Hi all,
> >>
> >> I'd like to revive this discussion. Since the last round I've re-aligned
> >> the proposal with current master, folded in all of the feedback,
> mirrored
> >> the proposal to a cwiki FLIP page (as Shengkai suggested), and opened a
> >> draft PR so the design can be reviewed against working code.
> >>
> >> Proposal doc (updated): link
> >> <
> https://docs.google.com/document/d/1WtAxp-jAVTLMOfWNldLCAoK137P0ZCMxR8hOZGcMxuc/edit
> >
> >> cwiki FLIP-527: link
> >> <
> https://cwiki.apache.org/confluence/spaces/FLINK/pages/353601981/FLIP-527+State+Schema+Evolution+for+RowData
> >
> >> Draft PR (implements the proposal; kept in draft until we converge):
> link
> >> <https://github.com/apache/flink/pull/28678>
> >>
> >> Here's what changed, mapped to the feedback:
> >>
> >> 1. Single object-level migration hook (Zakelly, Hangxiang). The design
> >> converged on one default method on TypeSerializerSnapshot: T
> >> migrate(TypeSerializerSnapshot<T> oldSerializerSnapshot, T value). It is
> >> invoked on the new snapshot and receives the old one (the
> >> oldSerializerSnapshot naming Hangxiang suggested), with an identity
> >> default. Because it works on the already-deserialized object rather
> than on
> >> raw bytes, the same method covers value, list-element, and map-value
> >> migration uniformly. This supersedes the earlier byte-level
> >> migrateState(in, out) form and removes the need for a separate
> >> migrateElement — which answers the migrateElement question you both
> raised
> >> (Zakelly, Hangxiang): the object hook handles list elements and map
> values
> >> with the same method, so no dedicated element hook is needed. It also
> >> addresses Hangxiang's point that a common interface shouldn't grow
> methods
> >> most implementations won't use, as they inherit the identity default and
> >> are entirely unaffected.
> >>
> >> 2. The SchemaEvolutionSerializer-on-the-compatibility-result alternative
> >> (Hangxiang). This is written up under Rejected Alternatives:
> >> TypeSerializerSchemaCompatibility is a result holder rather than an
> >> executor, and a migration is a distinct transform role (old serialized
> >> bytes -> new layout), so a default method on the snapshot is the
> smaller,
> >> more targeted change.
> >>
> >> 3. Field metadata (Shengkai). Adopted String[] fieldNames on
> >> RowDataSerializer instead of the full RowType. This approach is
> lightweight
> >> and exactly enough for name-based mapping.
> >>
> >> 4. The opt-in, and whether it is necessary (Shengkai, Gabor; and the
> >> question I left open last July). The feature is gated by
> >> table.exec.state.schema-evolution.enabled (default false), scoped under
> >> table.exec.state.* to signal it is Table/RowData-specific for now —
> which
> >> directly addresses Gabor's point that the config name should make clear
> it
> >> is RowData-specific rather than a generic state-evolution switch. On
> >> whether the opt-in is necessary: it is the fail-closed safety switch for
> >> exactly the case Shengkai raised, where a SQL change can silently shift
> >> operator-internal buffers (e.g. inserting SUM(d) before SUM(c)) even
> when
> >> field names appear to match. With the option off, serializers are built
> >> name-less and behavior is byte-for-byte as today; enabling it is a
> >> deliberate per-job confirmation that the change preserves state mapping.
> >>
> >> 5. Scope and a concrete example (Shengkai, Hongshun). The FLIP now leads
> >> with the primary case Shengkai steered toward — the SQL is unchanged and
> >> the input schema evolves backward-compatibly — with a worked end-to-end
> >> example: a Kafka fact stream joined with a CDC dimension whose nested
> >> profile ROW gains a nullable field. Because the join buffers the whole
> >> dimension row in keyed state, the evolved ROW actually reaches state (a
> >> leaf projection like metadata.userId would be column-pruned and would
> not),
> >> so the savepoint restore exercises exactly this feature. That also
> serves
> >> as the connector example Hongshun asked for.
> >>
> >> On Shengkai's broader concern (July 23) that the feature may not be
> >> accessible to users because few understand the SQL operator state
> >> structure: the reframing above is my attempt to address it — by leading
> >> with the SQL-unchanged case and a concrete join example rather than the
> >> operator-internal view. Shengkai, I'd especially welcome your read on
> >> whether this framing now makes the feature's applicability clear.
> >>
> >> Two boundaries I've made explicit this round:
> >> - RocksDB is the initial target backend; the ForSt sync backend can
> >> follow, and ForSt async is out of scope.
> >> - RowData nested below a composite serializer (List or Tuple, e.g.
> >> interval- and outer-join buffers) is out of scope for now and fails
> closed
> >> (rejected on restore, never mis-migrated); propagating migration through
> >> composite serializers is a planned follow-up.
> >>
> >> The proposal doc and cwiki have the full details, worked examples, the
> >> V3->V4 snapshot compatibility story, and rejected alternatives.
> Feedback is
> >> very welcome. If the direction looks good after this round, I'll start a
> >> VOTE.
> >>
> >> Thanks again for all the input,
> >> Weiqing
> >>
> >>
> >> On Tue, Aug 19, 2025 at 7:34 AM Gabor Somogyi <[email protected]
> >
> >> wrote:
> >>
> >>> Hi Weiqing,
> >>>
> >>> I've just read through the whole FLIP and +1 on the direction.
> >>>
> >>> I've a comment apart from the other pending items. Namely the
> >>> configuration is
> >>> `state.schema-evolution.enable` which implied to me that it's a generic
> >>> state evolution
> >>> feature but it's limited to Row data. Maybe we can mark that it's Row
> >>> data specific.
> >>> I'm pretty sure that we're going to add further types but not all.
> >>>
> >>> BR,
> >>> G
> >>>
> >>> On 2025/04/26 05:45:32 Weiqing Yang wrote:
> >>>> Hi all,
> >>>>
> >>>> I’d like to initiate a discussion about enhancing state schema
> evolution
> >>>> support for RowData in Flink.
> >>>>
> >>>> *Motivation*
> >>>>
> >>>> Flink applications frequently need to evolve their state schema as
> >>> business
> >>>> requirements change. Currently, when users update a Table API or SQL
> job
> >>>> with schema changes involving RowData types (particularly nested
> >>>> structures), they encounter serialization compatibility errors during
> >>> state
> >>>> restoration, causing job failures.The issue occurs because existing
> >>> state
> >>>> migration mechanisms don't properly handle RowData types during schema
> >>>> evolution, preventing users from making backward-compatible changes
> >>> like:
> >>>>
> >>>>   -
> >>>>
> >>>>   Adding nullable fields to existing structures
> >>>>   -
> >>>>
> >>>>   Reordering fields within a row while preserving field names
> >>>>   -
> >>>>
> >>>>   Evolving nested row structures
> >>>>
> >>>> This limitation impacts production applications using Flink's Table
> >>> API, as
> >>>> the RowData type is central to this interface. Users are forced to
> >>> choose
> >>>> between maintaining outdated schemas or reprocessing all state data
> when
> >>>> schema changes are required.
> >>>>
> >>>> Here’s the proposal document: Link
> >>>> <
> >>>
> https://docs.google.com/document/d/1WtAxp-jAVTLMOfWNldLCAoK137P0ZCMxR8hOZGcMxuc/edit?tab=t.0
> >>>>
> >>>> Your feedback and ideas are welcome to refine this feature.
> >>>>
> >>>> Thanks,
> >>>> Weiqing
> >>>>
> >>>
> >>
>
>

Reply via email to